feat(generation): worker孤儿任务自动恢复 + 429限流结构化提示 (#1677) #1710

Merged
auto-approve-bot merged 1 commits from feature/task-fault-tolerance-orphan-recovery into develop 2026-09-05 11:37:51 +08:00
Owner

背景

接 Issue #1677/#1702。问题:worker 容器部署重启时,正在跑的 preview/正式生成任务永久卡在 pending/running(实测 3 个预览任务卡 80% 超 10 小时);且旧任务占限流名额导致新请求 429,前端误报"创建失败"。

改动

1. worker 孤儿任务自动恢复(定时巡检,不再只靠启动清理)

改动前 改动后
running 孤儿清理 worker_ready 启动时执行一次,worker 不重启但任务卡死(OOM/上传挂起)无人管 beat 每 5 分钟巡检一次 + 启动时清理
running 超时阈值 10 分钟 20 分钟(任务硬超时 11 分钟,正常任务绝不可能超过,不误杀)
pending 超时阈值 30 分钟(僵尸长期占限流名额) 15 分钟(满队列 20 个 × 均 2 分钟 / 并发 4 ≈ 10 分钟,留余量)
pending 巡检周期 10 分钟 5 分钟
失败原因 文本 error_info.error_type = WorkerInterrupted / PendingTimeout + 完成时间
  • 巡检任务:worker.cleanup_stale_running_tasks(cleanup.py),同时清 generation_tasks 和 jobs 表孤儿
  • pending 被清理后限流名额立即释放(限流计数只数 pending 状态),新请求不再被 429 误伤

2. 429/503 限流结构化提示(前端区分"排队"与"创建失败")

429/503 响应体从纯文本改为结构化 detail(预览/正式生成/确认生成/重试共 7 处统一):

// 429 用户排队(用户自己的任务在跑/排队)
{
  "code": "USER_QUEUE_FULL",
  "message": "您有 3 个任务正在排队、2 个正在渲染,同一时间最多提交 3 个任务。请等待约 2 分钟后再提交",
  "queued_count": 3,
  "running_count": 2,
  "queue_ahead": 3,
  "estimated_wait_seconds": 120,
  "limit": 3
}

// 503 系统繁忙
{
  "code": "SYSTEM_QUEUE_FULL",
  "message": "系统繁忙:当前 20 个任务排队中、4 个渲染中,预计等待约 10 分钟,请稍后再试",
  "queued_count": 20, "running_count": 4, "queue_ahead": 20,
  "estimated_wait_seconds": 600, "limit": 20
}
  • code 是前端分支依据USER_QUEUE_FULL → 展示"排队中/稍后自动继续";SYSTEM_QUEUE_FULL → "系统繁忙稍后重试";创建失败(500)维持原逻辑
  • 等待预估:ceil(排队数 / worker 并发 4) × 最近 20 条已完成任务平均耗时(无历史数据默认 120 秒/条)
  • 兼容:新仓储方法全部鸭子类型防御(getattr),旧 mock/调用方零破坏

3. 仓储新增方法

  • count_running_by_user(user_id) / count_running_total():渲染中任务计数
  • estimate_avg_duration_seconds():最近完成任务平均耗时(Python 侧算时间差,兼容 SQLite/PostgreSQL)

测试

新增 tests/unit/test_task_fault_tolerance_1709.py(15 个):

  • 中断 running 任务(含 25 分钟无更新的预览任务)被重置 failed + 原因正确
  • 卡 pending 15 分钟任务被清理且 pending 计数归零(名额释放)
  • 3 个卡 10 小时的孤儿批量恢复(工单实测场景)
  • 正常 running(5 分钟前更新)/ 正常 pending(3 分钟)不被误杀
  • 限流计数:running 计数含预览任务、按用户隔离
  • build_rate_limit_detail:USER_QUEUE_FULL / SYSTEM_QUEUE_FULL 字段完整、不泄露内部 ID
  • 等待预估:并发换算(8 排队 / 4 并发 = 2 批)、仓储异常/旧仓储优雅降级

全量 14241 passed;black/isort/ruff/mypy 通过。

前端对接说明

429/503 现在 detail 是对象(不再是字符串),判断方式:

  • detail.code === "USER_QUEUE_FULL"(HTTP 429):提示排队,可展示 detail.message 或自行用 estimated_wait_seconds 拼"预计等待 X 分钟"
  • detail.code === "SYSTEM_QUEUE_FULL"(HTTP 503):系统繁忙
  • 其他 500/4xx 维持原错误处理
## 背景 接 Issue #1677/#1702。**问题**:worker 容器部署重启时,正在跑的 preview/正式生成任务永久卡在 pending/running(实测 3 个预览任务卡 80% 超 10 小时);且旧任务占限流名额导致新请求 429,前端误报"创建失败"。 ## 改动 ### 1. worker 孤儿任务自动恢复(定时巡检,不再只靠启动清理) | 项 | 改动前 | 改动后 | |---|---|---| | running 孤儿清理 | 仅 `worker_ready` 启动时执行一次,worker 不重启但任务卡死(OOM/上传挂起)无人管 | beat 每 **5 分钟**巡检一次 + 启动时清理 | | running 超时阈值 | 10 分钟 | **20 分钟**(任务硬超时 11 分钟,正常任务绝不可能超过,不误杀) | | pending 超时阈值 | 30 分钟(僵尸长期占限流名额) | **15 分钟**(满队列 20 个 × 均 2 分钟 / 并发 4 ≈ 10 分钟,留余量) | | pending 巡检周期 | 10 分钟 | 5 分钟 | | 失败原因 | 文本 | `error_info.error_type = WorkerInterrupted / PendingTimeout` + 完成时间 | - 巡检任务:`worker.cleanup_stale_running_tasks`(cleanup.py),同时清 generation_tasks 和 jobs 表孤儿 - pending 被清理后**限流名额立即释放**(限流计数只数 pending 状态),新请求不再被 429 误伤 ### 2. 429/503 限流结构化提示(前端区分"排队"与"创建失败") 429/503 响应体从纯文本改为结构化 `detail`(预览/正式生成/确认生成/重试共 7 处统一): ```json // 429 用户排队(用户自己的任务在跑/排队) { "code": "USER_QUEUE_FULL", "message": "您有 3 个任务正在排队、2 个正在渲染,同一时间最多提交 3 个任务。请等待约 2 分钟后再提交", "queued_count": 3, "running_count": 2, "queue_ahead": 3, "estimated_wait_seconds": 120, "limit": 3 } // 503 系统繁忙 { "code": "SYSTEM_QUEUE_FULL", "message": "系统繁忙:当前 20 个任务排队中、4 个渲染中,预计等待约 10 分钟,请稍后再试", "queued_count": 20, "running_count": 4, "queue_ahead": 20, "estimated_wait_seconds": 600, "limit": 20 } ``` - **`code` 是前端分支依据**:`USER_QUEUE_FULL` → 展示"排队中/稍后自动继续";`SYSTEM_QUEUE_FULL` → "系统繁忙稍后重试";创建失败(500)维持原逻辑 - 等待预估:`ceil(排队数 / worker 并发 4) × 最近 20 条已完成任务平均耗时`(无历史数据默认 120 秒/条) - 兼容:新仓储方法全部鸭子类型防御(`getattr`),旧 mock/调用方零破坏 ### 3. 仓储新增方法 - `count_running_by_user(user_id)` / `count_running_total()`:渲染中任务计数 - `estimate_avg_duration_seconds()`:最近完成任务平均耗时(Python 侧算时间差,兼容 SQLite/PostgreSQL) ## 测试 新增 `tests/unit/test_task_fault_tolerance_1709.py`(15 个): - 中断 running 任务(含 25 分钟无更新的预览任务)被重置 failed + 原因正确 - 卡 pending 15 分钟任务被清理且 pending 计数归零(名额释放) - 3 个卡 10 小时的孤儿批量恢复(工单实测场景) - 正常 running(5 分钟前更新)/ 正常 pending(3 分钟)不被误杀 - 限流计数:running 计数含预览任务、按用户隔离 - `build_rate_limit_detail`:USER_QUEUE_FULL / SYSTEM_QUEUE_FULL 字段完整、不泄露内部 ID - 等待预估:并发换算(8 排队 / 4 并发 = 2 批)、仓储异常/旧仓储优雅降级 全量 **14241 passed**;black/isort/ruff/mypy 通过。 ## 前端对接说明 429/503 现在 `detail` 是对象(不再是字符串),判断方式: - `detail.code === "USER_QUEUE_FULL"`(HTTP 429):提示排队,可展示 `detail.message` 或自行用 `estimated_wait_seconds` 拼"预计等待 X 分钟" - `detail.code === "SYSTEM_QUEUE_FULL"`(HTTP 503):系统繁忙 - 其他 500/4xx 维持原错误处理
xiaoxia added 1 commit 2026-09-05 11:26:44 +08:00
feat(generation): worker孤儿任务自动恢复+429限流结构化提示
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 26s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 36s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m25s
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m41s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m52s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m15s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m9s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 5m16s
AI Code Review / AI Code Review (pull_request) Successful in 6m21s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 10m15s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 4s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 7m53s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 8s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 24s
ce943faaeb
问题:
- worker容器部署重启时,正在跑的preview/正式生成任务永久卡在pending/running
  (实测3个预览任务卡80%超10小时),孤儿任务占限流名额导致新请求429误报创建失败
- 429返回纯文本,前端无法区分「限流排队」和「创建失败」

改动:
1. worker孤儿任务巡检(apps/worker/worker_app/tasks/)
   - running孤儿清理从仅worker_ready启动时执行,扩展为beat每5分钟定时巡检
     (cleanup_stale_running_tasks),容器重启/进程OOM卡死持续兜底
   - 阈值: running 10→20分钟(任务硬超时11分钟,20分钟绝不误杀);
     pending 30→15分钟(满队列消化约10分钟,留余量)
   - pending beat巡检10→5分钟,孤儿占位释放更快
   - 失败原因结构化: error_type=WorkerInterrupted/PendingTimeout
2. 429/503限流结构化提示(app/core/task_enqueue.py)
   - build_rate_limit_detail(): code/message/queued_count/running_count/
     queue_ahead/estimated_wait_seconds/limit
   - code: USER_QUEUE_FULL(429排队等待) / SYSTEM_QUEUE_FULL(503系统繁忙)
   - 等待预估: ceil(排队数/worker并发4) × 最近20条完成任务平均耗时
   - 预览/正式生成/确认生成/重试共7处限流响应全部结构化
   - 新仓储方法鸭子类型防御,旧mock/调用方零破坏
3. 仓储新增: count_running_by_user/count_running_total/estimate_avg_duration_seconds

测试: 新增15个单测(中断running重置、stale pending重置释放名额、3孤儿批量恢复、
限流计数含预览任务、结构化detail字段、等待预估并发计算、旧仓储优雅降级);
全量14241 passed; black/isort/ruff/mypy通过

🚀 预览环境已部署

项目 详情
PR号 #1710
预览链接 https://pr-1710.preview.xiaoxiajianji.com
API环境 staging

💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。

🔄 每次提交新代码后预览环境会自动更新。

🗑️ PR 关闭或合并后,预览环境会自动清理。

🚀 **预览环境已部署** | 项目 | 详情 | |------|------| | PR号 | #1710 | | 预览链接 | [https://pr-1710.preview.xiaoxiajianji.com](https://pr-1710.preview.xiaoxiajianji.com) | | API环境 | staging | > 💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。 > > 🔄 每次提交新代码后预览环境会自动更新。 > > 🗑️ PR 关闭或合并后,预览环境会自动清理。
Collaborator

【阻塞级判定】

  • 是否存在阻塞级问题:否
  • 阻塞级问题数量:0 个

📊 审查概览

  • 整体评价:通过
  • 建议级问题数量:3 个

🔴 阻塞级问题(必须修复)

💡 改进建议(不阻塞合并)

  1. [apps/api/app/core/task_enqueue.py: build_rate_limit_detail] 性能风险

    • 具体内容:build_rate_limit_detail 函数在错误处理路径中调用了 estimate_avg_duration_seconds,该方法会执行数据库查询 (ORDER BY ... LIMIT 20)。在系统高负载(大量触发 429/503)的情况下,频繁的数据库查询可能加剧数据库压力,导致响应变慢甚至雪崩。建议将平均耗时的计算结果缓存(如 Redis 或内存缓存,有效期 1-5 分钟),避免在每次限流错误时都查询数据库。
  2. [apps/api/app/api/routes/generation_tasks.py: create_generation_task] 逻辑一致性

    • 具体内容:在 create_generation_task 的循环中(diff 579-588行),当发生 UserPendingLimitExceededGlobalQueueFull 时,代码直接 break 跳出循环,但并未将当前 created_tasks 中最后一个任务标记为失败。这会导致该任务在数据库中处于 pendingcreated 状态但未入队,成为孤儿任务(依赖后台清理任务回收)。而在 generation_preview.py 中(diff 515-528行),同样的场景调用了 _mark_task_failed。建议统一行为,在 break 前显式标记任务失败,确保数据一致性。
  3. [apps/api/app/api/routes/generation_tasks.py: retry_generation_task] 代码设计

    • 具体内容:在 retry_generation_task 中(diff 804-815行),为了复用 build_rate_limit_detail 函数,代码手动实例化了 UserPendingLimitExceededGlobalQueueFull 异常对象仅作为数据容器传递。这种“用异常传数据”的模式不够直观。建议重构 build_rate_limit_detail,使其接受一个包含 pending_count, limit 等字段的数据类或字典,而不是强依赖异常类实例,以提高代码可读性和解耦。

良好实践

  • generation_preview.py 的循环处理中,对入队失败的任务显式调用 _mark_task_failed,有效避免了孤儿任务的产生。
  • build_rate_limit_detail 使用 getattr 检查仓储层方法是否存在,提供了良好的向后兼容性和防御性编程。
  • cleanup_orphan_tasks 等清理任务中使用了 try...finally 确保数据库 Session 正确关闭,避免了连接泄露。
  • estimate_avg_duration_seconds 中增加了对负数时长的过滤(total_seconds() > 0),防止时钟回拨导致的数据异常。

🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 03:32:59 | 模型:

### 【阻塞级判定】 - 是否存在阻塞级问题:否 - 阻塞级问题数量:0 个 ### 📊 审查概览 - 整体评价:通过 - 建议级问题数量:3 个 ### 🔴 阻塞级问题(必须修复) 无 ### 💡 改进建议(不阻塞合并) 1. **[apps/api/app/core/task_enqueue.py: build_rate_limit_detail] 性能风险** - 具体内容:`build_rate_limit_detail` 函数在错误处理路径中调用了 `estimate_avg_duration_seconds`,该方法会执行数据库查询 (`ORDER BY ... LIMIT 20`)。在系统高负载(大量触发 429/503)的情况下,频繁的数据库查询可能加剧数据库压力,导致响应变慢甚至雪崩。建议将平均耗时的计算结果缓存(如 Redis 或内存缓存,有效期 1-5 分钟),避免在每次限流错误时都查询数据库。 2. **[apps/api/app/api/routes/generation_tasks.py: create_generation_task] 逻辑一致性** - 具体内容:在 `create_generation_task` 的循环中(diff 579-588行),当发生 `UserPendingLimitExceeded` 或 `GlobalQueueFull` 时,代码直接 `break` 跳出循环,但并未将当前 `created_tasks` 中最后一个任务标记为失败。这会导致该任务在数据库中处于 `pending` 或 `created` 状态但未入队,成为孤儿任务(依赖后台清理任务回收)。而在 `generation_preview.py` 中(diff 515-528行),同样的场景调用了 `_mark_task_failed`。建议统一行为,在 `break` 前显式标记任务失败,确保数据一致性。 3. **[apps/api/app/api/routes/generation_tasks.py: retry_generation_task] 代码设计** - 具体内容:在 `retry_generation_task` 中(diff 804-815行),为了复用 `build_rate_limit_detail` 函数,代码手动实例化了 `UserPendingLimitExceeded` 和 `GlobalQueueFull` 异常对象仅作为数据容器传递。这种“用异常传数据”的模式不够直观。建议重构 `build_rate_limit_detail`,使其接受一个包含 `pending_count`, `limit` 等字段的数据类或字典,而不是强依赖异常类实例,以提高代码可读性和解耦。 ### ✅ 良好实践 - 在 `generation_preview.py` 的循环处理中,对入队失败的任务显式调用 `_mark_task_failed`,有效避免了孤儿任务的产生。 - `build_rate_limit_detail` 使用 `getattr` 检查仓储层方法是否存在,提供了良好的向后兼容性和防御性编程。 - `cleanup_orphan_tasks` 等清理任务中使用了 `try...finally` 确保数据库 Session 正确关闭,避免了连接泄露。 - 在 `estimate_avg_duration_seconds` 中增加了对负数时长的过滤(`total_seconds() > 0`),防止时钟回拨导致的数据异常。 --- <sub>🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 03:32:59 | 模型: </sub> <!-- AI_CODE_REVIEW_AUTO_COMMENT -->
auto-approve-bot merged commit 21c26b5b26 into develop 2026-09-05 11:37:51 +08:00
auto-approve-bot deleted branch feature/task-fault-tolerance-orphan-recovery 2026-09-05 11:37:51 +08:00

🗑️ 预览环境已清理

PR #1710 已关闭或合并,对应的预览环境已被清理。

如有需要,可以重新打开 PR 来重新生成预览环境。

🗑️ **预览环境已清理** PR #1710 已关闭或合并,对应的预览环境已被清理。 > 如有需要,可以重新打开 PR 来重新生成预览环境。
Sign in to join this conversation.