feat(cleanup): pending 任务自动清理,防止队列占满 #1463
Reference in New Issue
Block a user
Delete Branch "feat/pending-timeout-cleanup"
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?
背景
全局任务队列上限 20 个 pending,一旦有任务卡在 pending(Worker 没拉取、异常中断等),队列满了后新用户无法创建任务(返回 503)。
方案
在 Celery Beat 定时任务中新增定期清理任务:
改动
generation_task_repository.py:新增cleanup_stale_pending()方法status=pending且created_at < NOW() - timeout_minutes的 tasksstatus=failed,error_message="pending timeout: auto cleanup",completed_at=NOW()error_info含PendingTimeout类型和时间戳_startup.py:新增cleanup_stale_pending_tasks()函数PENDING_TASK_TIMEOUT_MINUTES = 30常量cleanup_all_stale_tasks()也清理 pendingcleanup.py(新文件):Celery 定时任务@shared_task(name="worker.cleanup_stale_pending_tasks")_startup.cleanup_stale_pending_tasks()celery_app.py:添加 beat_scheduleexpires=300防止任务堆积entrypoint-worker.sh:添加-B标志测试:6 个单元测试覆盖正常/超时/批量/不干扰 running 等场景
阈值配置
PENDING_TASK_TIMEOUT_MINUTES = 30:默认 30 分钟,可调整🚀 预览环境已部署
CI全绿,自动审批通过。
CI全绿,自动审批通过。
【阻塞级判定】
📊 审查概览
🔴 阻塞级问题(必须修复)
entrypoint-worker.sh中添加-B参数启动了内嵌的 Beat 调度器。如果 Worker 部署了多个实例(例如 K8sreplicas > 1),每个实例都会独立运行 Beat 调度器,导致定时任务cleanup-stale-pending-tasks被并发执行多次。虽然数据库更新操作是幂等的(状态变为 FAILED),但这会导致重复的数据库写入压力、重复的日志记录以及资源浪费。-B。建议部署独立的 Beat 服务(如celery -A worker_app.celery_app beat),或者确保基础设施层面保证该 Worker 永远只运行单实例。若仅为单实例部署,请在文档或部署配置中明确标注此限制。💡 改进建议(不阻塞合并)
[packages/adapters/sqlalchemy_impl/generation_task_repository.py: 337] 批量更新性能优化
synchronize_session="fetch"会在更新后执行 SELECT 语句来同步会话中的对象状态。如果pending状态的任务数量较多(尽管有上限,但在极端情况下),这会产生额外的查询开销。由于此方法通常在独立的清理任务中运行,且后续不依赖会话中的对象状态,建议使用synchronize_session=False以提升性能。[apps/worker/worker_app/celery_app.py: 28] 任务过期时间可能导致清理延迟
options: {"expires": 300}表示如果任务在队列中等待超过 5 分钟未被执行,将被丢弃。考虑到调度间隔为 10 分钟,如果 Worker 繁忙导致任务堆积,清理任务可能会被跳过,导致实际清理间隔延长至 20 分钟。建议评估是否需要expires参数,或者适当调大该值以防止清理频率低于预期。✅ 良好实践
update方法,避免了循环查询,效率较高。✅ 格式检查通过 | ❌ 逻辑审查需修改 | ⚠️ 建议关注性能
🤖 由 AI 代码审查机器人自动生成 | 2026-08-23 04:41:59 | 模型:
🗑️ 预览环境已清理
PR #1463 已关闭或合并,对应的预览环境已被清理。