feat(generation): worker孤儿任务自动恢复 + 429限流结构化提示 (#1677) #1710
Reference in New Issue
Block a user
Delete Branch "feature/task-fault-tolerance-orphan-recovery"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
背景
接 Issue #1677/#1702。问题:worker 容器部署重启时,正在跑的 preview/正式生成任务永久卡在 pending/running(实测 3 个预览任务卡 80% 超 10 小时);且旧任务占限流名额导致新请求 429,前端误报"创建失败"。
改动
1. worker 孤儿任务自动恢复(定时巡检,不再只靠启动清理)
worker_ready启动时执行一次,worker 不重启但任务卡死(OOM/上传挂起)无人管error_info.error_type = WorkerInterrupted / PendingTimeout+ 完成时间worker.cleanup_stale_running_tasks(cleanup.py),同时清 generation_tasks 和 jobs 表孤儿2. 429/503 限流结构化提示(前端区分"排队"与"创建失败")
429/503 响应体从纯文本改为结构化
detail(预览/正式生成/确认生成/重试共 7 处统一):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 个):build_rate_limit_detail:USER_QUEUE_FULL / SYSTEM_QUEUE_FULL 字段完整、不泄露内部 ID全量 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):系统繁忙问题: - 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通过🚀 预览环境已部署
【阻塞级判定】
📊 审查概览
🔴 阻塞级问题(必须修复)
无
💡 改进建议(不阻塞合并)
[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 分钟),避免在每次限流错误时都查询数据库。[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前显式标记任务失败,确保数据一致性。[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),防止时钟回拨导致的数据异常。🤖 由 AI 代码审查机器人自动生成 | 2026-09-05 03:32:59 | 模型:
🗑️ 预览环境已清理
PR #1710 已关闭或合并,对应的预览环境已被清理。