feat(cleanup): pending 任务自动清理,防止队列占满 #1463

Merged
xiaoxia merged 4 commits from feat/pending-timeout-cleanup into develop 2026-08-23 12:55:40 +08:00
Owner

背景

全局任务队列上限 20 个 pending,一旦有任务卡在 pending(Worker 没拉取、异常中断等),队列满了后新用户无法创建任务(返回 503)。

方案

在 Celery Beat 定时任务中新增定期清理任务:

改动

  1. generation_task_repository.py:新增 cleanup_stale_pending() 方法

    • 查找 status=pendingcreated_at < NOW() - timeout_minutes 的 tasks
    • 批量更新为 status=failed, error_message="pending timeout: auto cleanup", completed_at=NOW()
    • 填充 error_infoPendingTimeout 类型和时间戳
  2. _startup.py:新增 cleanup_stale_pending_tasks() 函数

    • 定义 PENDING_TASK_TIMEOUT_MINUTES = 30 常量
    • Worker 启动时 cleanup_all_stale_tasks() 也清理 pending
  3. cleanup.py(新文件):Celery 定时任务

    • @shared_task(name="worker.cleanup_stale_pending_tasks")
    • 委托给 _startup.cleanup_stale_pending_tasks()
  4. celery_app.py:添加 beat_schedule

    • 每 600 秒(10 分钟)执行一次
    • expires=300 防止任务堆积
  5. entrypoint-worker.sh:添加 -B 标志

    • Worker 进程内嵌 celery beat(单实例部署足够安全)
  6. 测试:6 个单元测试覆盖正常/超时/批量/不干扰 running 等场景

阈值配置

  • PENDING_TASK_TIMEOUT_MINUTES = 30:默认 30 分钟,可调整
  • Beat 间隔 10 分钟,expires 5 分钟
## 背景 全局任务队列上限 20 个 pending,一旦有任务卡在 pending(Worker 没拉取、异常中断等),队列满了后新用户无法创建任务(返回 503)。 ## 方案 在 Celery Beat 定时任务中新增定期清理任务: ### 改动 1. **`generation_task_repository.py`**:新增 `cleanup_stale_pending()` 方法 - 查找 `status=pending` 且 `created_at < NOW() - timeout_minutes` 的 tasks - 批量更新为 `status=failed`, `error_message="pending timeout: auto cleanup"`, `completed_at=NOW()` - 填充 `error_info` 含 `PendingTimeout` 类型和时间戳 2. **`_startup.py`**:新增 `cleanup_stale_pending_tasks()` 函数 - 定义 `PENDING_TASK_TIMEOUT_MINUTES = 30` 常量 - Worker 启动时 `cleanup_all_stale_tasks()` 也清理 pending 3. **`cleanup.py`(新文件)**:Celery 定时任务 - `@shared_task(name="worker.cleanup_stale_pending_tasks")` - 委托给 `_startup.cleanup_stale_pending_tasks()` 4. **`celery_app.py`**:添加 beat_schedule - 每 600 秒(10 分钟)执行一次 - `expires=300` 防止任务堆积 5. **`entrypoint-worker.sh`**:添加 `-B` 标志 - Worker 进程内嵌 celery beat(单实例部署足够安全) 6. **测试**:6 个单元测试覆盖正常/超时/批量/不干扰 running 等场景 ## 阈值配置 - `PENDING_TASK_TIMEOUT_MINUTES = 30`:默认 30 分钟,可调整 - Beat 间隔 10 分钟,expires 5 分钟
xiaoxia added 1 commit 2026-08-23 12:22:09 +08:00
feat(cleanup): pending task auto cleanup via Celery Beat
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 / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
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 / Check if frontend-only change (pull_request) Successful in 42s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Validate - Migration (alembic) (pull_request) Successful in 1m35s
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Validate - Type Check (mypy) (pull_request) Successful in 1m37s
AI Code Review / AI Code Review (pull_request) Failing after 1m38s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m55s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 33s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m42s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m23s
CI/CD Pipeline / Validate - Code Quality (pull_request) Has been cancelled
CI/CD Pipeline / Unit Tests (pull_request) Has been cancelled
CI/CD Pipeline / Integration Tests (pull_request) Has been cancelled
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
PR Automation / Auto Approve on CI Green (pull_request) Has been cancelled
d05279b410
- generation_task_repository: add cleanup_stale_pending() method
  marks status=pending + created_at > timeout as failed
  with error_message="pending timeout: auto cleanup"
- _startup.py: add cleanup_stale_pending_tasks() function
  with PENDING_TASK_TIMEOUT_MINUTES = 30 constant
- cleanup.py: new Celery scheduled task worker.cleanup_stale_pending_tasks
- celery_app.py: add beat_schedule, runs every 10 minutes (600s)
- entrypoint-worker.sh: add -B flag for embedded celery beat
- 6 unit tests for cleanup_stale_pending

🚀 预览环境已部署

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

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

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

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

🚀 **预览环境已部署** | 项目 | 详情 | |------|------| | PR号 | #1463 | | 预览链接 | [https://pr-1463.preview.xiaoxiajianji.com](https://pr-1463.preview.xiaoxiajianji.com) | | API环境 | staging | > 💡 预览环境使用 staging API 数据,请勿在预览环境中操作重要数据。 > > 🔄 每次提交新代码后预览环境会自动更新。 > > 🗑️ PR 关闭或合并后,预览环境会自动清理。
auto-approve-bot added 1 commit 2026-08-23 12:25:01 +08:00
style: auto-format with black + isort + prettier [skip ci-format-check]
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 / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
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 / Check if frontend-only change (pull_request) Successful in 46s
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 / Validate - Type Check (mypy) (pull_request) Successful in 1m37s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 45s
AI Code Review / AI Code Review (pull_request) Failing after 1m40s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m44s
CI/CD Pipeline / Validate - Migration (alembic) (pull_request) Successful in 1m50s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m28s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m47s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 2m35s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m30s
CI/CD Pipeline / Validate - Code Quality (pull_request) Successful in 4m55s
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 / Integration Tests (pull_request) Successful in 1m37s
CI/CD Pipeline / CI Gate (pull_request) Successful in 8s
d3a9fff0c2
auto-approve-bot approved these changes 2026-08-23 12:30:31 +08:00
auto-approve-bot left a comment
Collaborator

CI全绿,自动审批通过。

CI全绿,自动审批通过。
auto-approve-bot approved these changes 2026-08-23 12:30:31 +08:00
auto-approve-bot left a comment
Collaborator

CI全绿,自动审批通过。

CI全绿,自动审批通过。
xiaoxia added 1 commit 2026-08-23 12:39:53 +08:00
fix: address AI review - session leak & bulk update for pending cleanup
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 / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
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 / Check if frontend-only change (pull_request) Successful in 42s
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 Worker Image (pull_request) Successful in 47s
CI/CD Pipeline / Validate - Type Check (mypy) (pull_request) Successful in 1m40s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m44s
CI/CD Pipeline / Validate - Migration (alembic) (pull_request) Successful in 1m50s
AI Code Review / AI Code Review (pull_request) Failing after 2m8s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m16s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m41s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 2m32s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m28s
CI/CD Pipeline / Validate - Code Quality (pull_request) Has been cancelled
CI/CD Pipeline / Integration Tests (pull_request) Has been cancelled
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
1712c1693a
- _startup.py: move session.close() to finally block to prevent DB connection leak
- generation_task_repository.py: replace .all()+loop with bulk .update() to avoid OOM
  when large number of pending tasks accumulate
Collaborator

【阻塞级判定】

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

📊 审查概览

  • 整体评价:需修改
  • 建议级问题数量:2 个

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

  1. [infra/docker/entrypoint-worker.sh: 10] Celery Beat 多实例并发风险
    • 问题类型:逻辑bug / 数据损坏风险
    • 问题描述:在 entrypoint-worker.sh 中添加 -B 参数启动了内嵌的 Beat 调度器。如果 Worker 部署了多个实例(例如 K8s replicas > 1),每个实例都会独立运行 Beat 调度器,导致定时任务 cleanup-stale-pending-tasks 被并发执行多次。虽然数据库更新操作是幂等的(状态变为 FAILED),但这会导致重复的数据库写入压力、重复的日志记录以及资源浪费。
    • 修改建议:如果必须支持水平扩展,不应在 Worker 进程中混用 -B。建议部署独立的 Beat 服务(如 celery -A worker_app.celery_app beat),或者确保基础设施层面保证该 Worker 永远只运行单实例。若仅为单实例部署,请在文档或部署配置中明确标注此限制。

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

  1. [packages/adapters/sqlalchemy_impl/generation_task_repository.py: 337] 批量更新性能优化

    • 具体内容:synchronize_session="fetch" 会在更新后执行 SELECT 语句来同步会话中的对象状态。如果 pending 状态的任务数量较多(尽管有上限,但在极端情况下),这会产生额外的查询开销。由于此方法通常在独立的清理任务中运行,且后续不依赖会话中的对象状态,建议使用 synchronize_session=False 以提升性能。
  2. [apps/worker/worker_app/celery_app.py: 28] 任务过期时间可能导致清理延迟

    • 具体内容:配置中 options: {"expires": 300} 表示如果任务在队列中等待超过 5 分钟未被执行,将被丢弃。考虑到调度间隔为 10 分钟,如果 Worker 繁忙导致任务堆积,清理任务可能会被跳过,导致实际清理间隔延长至 20 分钟。建议评估是否需要 expires 参数,或者适当调大该值以防止清理频率低于预期。

良好实践

  • 单元测试覆盖全面,包含了边界条件(如时间阈值、状态过滤)和批量更新的测试。
  • 数据库操作使用了 ORM 的批量 update 方法,避免了循环查询,效率较高。
  • 异常处理得当,清理任务捕获了所有异常并记录日志,防止因单次清理失败导致调度器崩溃。

格式检查通过 | 逻辑审查需修改 | ⚠️ 建议关注性能


🤖 由 AI 代码审查机器人自动生成 | 2026-08-23 04:41:59 | 模型:

### 【阻塞级判定】 - 是否存在阻塞级问题:是 - 阻塞阻塞级问题数量:1 个 ### 📊 审查概览 - 整体评价:需修改 - 建议级问题数量:2 个 ### 🔴 阻塞级问题(必须修复) 1. **[infra/docker/entrypoint-worker.sh: 10] Celery Beat 多实例并发风险** - 问题类型:逻辑bug / 数据损坏风险 - 问题描述:在 `entrypoint-worker.sh` 中添加 `-B` 参数启动了内嵌的 Beat 调度器。如果 Worker 部署了多个实例(例如 K8s `replicas > 1`),每个实例都会独立运行 Beat 调度器,导致定时任务 `cleanup-stale-pending-tasks` 被并发执行多次。虽然数据库更新操作是幂等的(状态变为 FAILED),但这会导致重复的数据库写入压力、重复的日志记录以及资源浪费。 - 修改建议:如果必须支持水平扩展,不应在 Worker 进程中混用 `-B`。建议部署独立的 Beat 服务(如 `celery -A worker_app.celery_app beat`),或者确保基础设施层面保证该 Worker 永远只运行单实例。若仅为单实例部署,请在文档或部署配置中明确标注此限制。 ### 💡 改进建议(不阻塞合并) 1. **[packages/adapters/sqlalchemy_impl/generation_task_repository.py: 337] 批量更新性能优化** - 具体内容:`synchronize_session="fetch"` 会在更新后执行 SELECT 语句来同步会话中的对象状态。如果 `pending` 状态的任务数量较多(尽管有上限,但在极端情况下),这会产生额外的查询开销。由于此方法通常在独立的清理任务中运行,且后续不依赖会话中的对象状态,建议使用 `synchronize_session=False` 以提升性能。 2. **[apps/worker/worker_app/celery_app.py: 28] 任务过期时间可能导致清理延迟** - 具体内容:配置中 `options: {"expires": 300}` 表示如果任务在队列中等待超过 5 分钟未被执行,将被丢弃。考虑到调度间隔为 10 分钟,如果 Worker 繁忙导致任务堆积,清理任务可能会被跳过,导致实际清理间隔延长至 20 分钟。建议评估是否需要 `expires` 参数,或者适当调大该值以防止清理频率低于预期。 ### ✅ 良好实践 - 单元测试覆盖全面,包含了边界条件(如时间阈值、状态过滤)和批量更新的测试。 - 数据库操作使用了 ORM 的批量 `update` 方法,避免了循环查询,效率较高。 - 异常处理得当,清理任务捕获了所有异常并记录日志,防止因单次清理失败导致调度器崩溃。 --- ✅ 格式检查通过 | ❌ 逻辑审查需修改 | ⚠️ 建议关注性能 --- <sub>🤖 由 AI 代码审查机器人自动生成 | 2026-08-23 04:41:59 | 模型: </sub> <!-- AI_CODE_REVIEW_AUTO_COMMENT -->
xiaoxia added 1 commit 2026-08-23 12:44:57 +08:00
fix: document single-instance constraint for embedded Beat & optimize synchronize_session
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 / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
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 / Check if frontend-only change (pull_request) Successful in 38s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Validate - Type Check (mypy) (pull_request) Successful in 1m32s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 1m46s
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Validate - Migration (alembic) (pull_request) Successful in 2m2s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 29s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m22s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m48s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 2m40s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m33s
CI/CD Pipeline / Validate - Code Quality (pull_request) Successful in 4m11s
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 / Integration Tests (pull_request) Successful in 1m40s
CI/CD Pipeline / CI Gate (pull_request) Successful in 6s
AI Code Review / AI Code Review (pull_request) Successful in 6m48s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 1m18s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 1m24s
498b490bac
- entrypoint-worker.sh: prominent deployment constraint warning (replicas=1 only)
- generation_task_repository.py: synchronize_session=False for cleanup (no post-update session sync needed)
xiaoxia merged commit 5a5c653d2c into develop 2026-08-23 12:55:40 +08:00
xiaoxia deleted branch feat/pending-timeout-cleanup 2026-08-23 12:55:41 +08:00

🗑️ 预览环境已清理

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

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

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