feat: 任务队列限流防护 #220

Merged
xiaoxia merged 2 commits from feat/task-queue-rate-limit into develop 2026-07-11 16:29:08 +08:00
Owner

背景

staging 队列今天两次爆炸(135→595 个 pending),全是批量提交导致的。Worker 处理能力有限(并发1,单任务10+分钟),批量涌入就彻底堵死。

实现3层防护

1. 用户级限流(核心)

  • 每个用户同时 pending 的 generation 任务上限 3 个
  • 超限返回 429您的待处理任务过多,请等待完成后再提交

2. 全局限流(兜底)

  • 系统 pending 任务超过 20 个一律拒绝
  • 返回 503系统繁忙,请稍后再试

3. 双重检查

  • 创建前预检查(批量提交提前拦截)
  • 入队前兜底检查(防止并发竞态)

覆盖的入口

入口 文件 说明
一键生成 generation_tasks.py 批量创建 + 单任务重试
剪辑计划生成 edit_plans.py 剪辑计划渲染
任务中心重试 task_center.py 用户级重试 + 项目级重试

关键设计

  • 限流拒绝时标记 failed:不留下 pending 僵尸任务
  • 全局检查始终生效:即使不传 user_id 也受全局保护
  • Repository 降级兼容:没有 count_pending 方法时自动跳过,不阻塞旧代码
  • > 判断语义:入队函数里用 >(任务已创建),预检查用 >=(任务未创建),边界一致

测试

13 个单元测试,覆盖:

  • 正常提交通过 ✓
  • 用户超限被拒(429)✓
  • 全局超限被拒(503)✓
  • 全局优先于用户 ✓
  • 边界值(刚好等于上限不算超限)✓
  • 不传 user_id 跳过用户级但保留全局 ✓
  • repository 更新失败不崩溃 ✓

全量 1296 个测试通过

## 背景 staging 队列今天两次爆炸(135→595 个 pending),全是批量提交导致的。Worker 处理能力有限(并发1,单任务10+分钟),批量涌入就彻底堵死。 ## 实现3层防护 ### 1. 用户级限流(核心) - 每个用户同时 pending 的 generation 任务上限 **3 个** - 超限返回 **429**:`您的待处理任务过多,请等待完成后再提交` ### 2. 全局限流(兜底) - 系统 pending 任务超过 **20 个**一律拒绝 - 返回 **503**:`系统繁忙,请稍后再试` ### 3. 双重检查 - 创建前预检查(批量提交提前拦截) - 入队前兜底检查(防止并发竞态) ## 覆盖的入口 | 入口 | 文件 | 说明 | |------|------|------| | 一键生成 | generation_tasks.py | 批量创建 + 单任务重试 | | 剪辑计划生成 | edit_plans.py | 剪辑计划渲染 | | 任务中心重试 | task_center.py | 用户级重试 + 项目级重试 | ## 关键设计 - **限流拒绝时标记 failed**:不留下 pending 僵尸任务 - **全局检查始终生效**:即使不传 user_id 也受全局保护 - **Repository 降级兼容**:没有 count_pending 方法时自动跳过,不阻塞旧代码 - **`>` 判断语义**:入队函数里用 `>`(任务已创建),预检查用 `>=`(任务未创建),边界一致 ## 测试 13 个单元测试,覆盖: - 正常提交通过 ✓ - 用户超限被拒(429)✓ - 全局超限被拒(503)✓ - 全局优先于用户 ✓ - 边界值(刚好等于上限不算超限)✓ - 不传 user_id 跳过用户级但保留全局 ✓ - repository 更新失败不崩溃 ✓ 全量 1296 个测试通过 ✅
xiaoxia added 1 commit 2026-07-11 16:13:37 +08:00
实现3层队列防护,防止批量提交导致队列爆炸:

**用户级限流(核心)**
- 每个用户同时 pending 的 generation 任务上限 3 个
- 超限返回 429:您的待处理任务过多,请等待完成后再提交

**全局限流(兜底)**
- 系统 pending 任务超过 20 个一律拒绝
- 返回 503:系统繁忙,请稍后再试

**覆盖的入口**
- 一键生成(generation_tasks 批量创建 + 重试)
- 剪辑计划生成(edit_plans generate)
- 任务中心重试(task_center 用户级 + 项目级)

**实现细节**
- 预检查 + 入队前检查双重保障
- 限流拒绝时任务标记为 failed,避免 pending 僵尸
- repository 不支持计数时自动降级跳过(兼容旧代码)
- 全局检查始终生效,用户级检查需传 user_id

**新增**
- task_enqueue.py: UserPendingLimitExceeded / GlobalQueueFull 异常
- generation_task_repository: count_pending_by_user / count_pending_total
- 13个单元测试,覆盖正常/超限/全局/降级等场景
Author
Owner

PR #220 审查结论:⚠️ 需要修复 3 个问题(P1×2 + P2×1)

审查范围

任务队列限流防护:用户级 pending 限流(3个上限,429)、全局限流(20个上限,503)。

做得好的地方

  1. 4 处入口全覆盖:批量创建、generation 重试、任务中心用户级重试、任务中心项目级重试
  2. 限流检查在入队前:先检查再发送 Celery 任务,方向正确
  3. 超限后标记 failed:不会留下 pending 僵尸任务,和安全入队设计一致
  4. 单元测试覆盖完整:正常通过、用户超限、全局超限、不传 user_id 跳过用户级、repository.update 失败不崩溃 等场景都有
  5. 全局优先于用户级:系统级保护优先级高,设计合理

⚠️ P1-1:限流阈值硬编码 + 各入口判断条件不一致

问题

  • task_enqueue.pysafe_enqueue_generation_task 默认参数 user_pending_limit=3, global_pending_limit=20
  • generation_tasks.py 预检查:硬编码 user_pending + count > 3 / global_pending + count > 20
  • task_center.py 两个重试入口:硬编码 user_pending >= 3 / global_pending >= 20

三个不一致

  1. 硬编码:4 处入口的阈值都是写死的数字,没有统一来源。以后要改上限得改 5 个地方,容易漏
  2. 判断条件不统一safe_enqueue_generation_task>(超过才拒),task_center.py 预检查用 >=(等于就拒)。语义不一致会导致边界行为不同
  3. 批量创建预检查多减了 countpending_count - count 来显示当前数,逻辑绕

建议修复

  • 定义常量 DEFAULT_USER_PENDING_LIMIT = 3 / DEFAULT_GLOBAL_PENDING_LIMIT = 20
  • 所有入口引用同一常量
  • 统一判断语义:超过上限才拒绝(即 >),保持一致

⚠️ P1-2:批量创建的竞态条件

问题create_generation_task 的预检查在 for 循环外,但实际建任务和入队在循环内。如果两个请求同时到达,预检查都通过了,循环里建了 5 个 + 5 个 = 10 个 pending,远超用户 3 个上限。

虽然 safe_enqueue_generation_task 里有兜底检查,但有两个问题:

  1. 已经创建了一堆任务才发现超限:前面成功的 + 后面失败的,用户看到一半成功一半失败,体验不好
  2. DB 中留下了失败任务记录:这些被限流拒绝的任务会以 failed 状态留在库里,用户重试又会新增,积累垃圾数据

建议修复

  • 批量创建前预检查用 user_pending + count > limit,把本次要提交的数量算进去,从源头拒绝
  • 兜底逻辑保留(防并发),但预检查应该更严格

看了下代码,预检查确实用了 user_pending + count > 3,这个方向是对的。但有个 bug:> 3 意味着 4 个才超限,但用户上限是 3 个。如果用户已有 3 个 pending,再提交 1 个,3 + 1 > 3 = True 会被拒绝,这个是对的。但如果用户已有 2 个,再提交 2 个,2 + 2 > 3 = True 也会被拒,而如果改成逐个入队的话其实第 1 个是成功的。

这个不是 bug,是策略选择——批量提交要么全过要么全拒,比一半成功一半失败好。没问题。

真正的竞态问题:两个并发请求各提交 1 个,用户已有 2 个。预检查都通过(2+1=3 不大于3),然后各自入队,最终用户有 4 个 pending。这个在 DB 层面没有锁。

建议

  • 短期可以接受,靠兜底的 safe_enqueue_generation_task 里再查一次兜底(但也有竞态,只是窗口小一点)
  • 如果要彻底解决,需要用 DB 行锁或 Redis 原子计数。当前阶段可以先接受"尽力而为"的限流,注释说明一下就行

⚠️ P2:预检查逻辑重复

问题:task_center 里的 retry_task_by_idretry_project_task 有几乎一模一样的预检查代码(用户级+全局级),generation_tasks.py 里也有类似的预检查。总共 3 处重复。

建议

  • 把预检查抽成一个公共函数,比如 enforce_queue_limits(user_id, repo, extra_count=1),超限直接抛 HTTPException
  • 或者把限流异常统一在中间件/异常处理器里转成 HTTP 响应,路由层就不用 try/except 了

总结

级别 问题 影响
P1 限流阈值硬编码 + 判断条件不统一 维护成本高,边界行为不一致
P1 并发竞态条件 极端情况下可能突破限流
P2 预检查逻辑 3 处重复 维护成本高

P1 问题修复后即可合并。P2 是代码质量优化,可以后续重构。

## PR #220 审查结论:⚠️ 需要修复 3 个问题(P1×2 + P2×1) ### 审查范围 任务队列限流防护:用户级 pending 限流(3个上限,429)、全局限流(20个上限,503)。 ### ✅ 做得好的地方 1. **4 处入口全覆盖**:批量创建、generation 重试、任务中心用户级重试、任务中心项目级重试 2. **限流检查在入队前**:先检查再发送 Celery 任务,方向正确 3. **超限后标记 failed**:不会留下 pending 僵尸任务,和安全入队设计一致 4. **单元测试覆盖完整**:正常通过、用户超限、全局超限、不传 user_id 跳过用户级、repository.update 失败不崩溃 等场景都有 5. **全局优先于用户级**:系统级保护优先级高,设计合理 --- ### ⚠️ P1-1:限流阈值硬编码 + 各入口判断条件不一致 **问题**: - `task_enqueue.py` 中 `safe_enqueue_generation_task` 默认参数 `user_pending_limit=3, global_pending_limit=20` ✅ - `generation_tasks.py` 预检查:硬编码 `user_pending + count > 3` / `global_pending + count > 20` - `task_center.py` 两个重试入口:硬编码 `user_pending >= 3` / `global_pending >= 20` **三个不一致**: 1. **硬编码**:4 处入口的阈值都是写死的数字,没有统一来源。以后要改上限得改 5 个地方,容易漏 2. **判断条件不统一**:`safe_enqueue_generation_task` 用 `>`(超过才拒),`task_center.py` 预检查用 `>=`(等于就拒)。语义不一致会导致边界行为不同 3. **批量创建预检查多减了 count**:`pending_count - count` 来显示当前数,逻辑绕 **建议修复**: - 定义常量 `DEFAULT_USER_PENDING_LIMIT = 3` / `DEFAULT_GLOBAL_PENDING_LIMIT = 20` - 所有入口引用同一常量 - 统一判断语义:**超过上限才拒绝**(即 `>`),保持一致 --- ### ⚠️ P1-2:批量创建的竞态条件 **问题**:`create_generation_task` 的预检查在 for 循环外,但实际建任务和入队在循环内。如果两个请求同时到达,预检查都通过了,循环里建了 5 个 + 5 个 = 10 个 pending,远超用户 3 个上限。 虽然 `safe_enqueue_generation_task` 里有兜底检查,但有两个问题: 1. **已经创建了一堆任务才发现超限**:前面成功的 + 后面失败的,用户看到一半成功一半失败,体验不好 2. **DB 中留下了失败任务记录**:这些被限流拒绝的任务会以 failed 状态留在库里,用户重试又会新增,积累垃圾数据 **建议修复**: - 批量创建前预检查用 `user_pending + count > limit`,把本次要提交的数量算进去,从源头拒绝 - 兜底逻辑保留(防并发),但预检查应该更严格 看了下代码,预检查确实用了 `user_pending + count > 3`,这个方向是对的。但有个 bug:**`> 3` 意味着 4 个才超限,但用户上限是 3 个**。如果用户已有 3 个 pending,再提交 1 个,`3 + 1 > 3 = True` 会被拒绝,这个是对的。但如果用户已有 2 个,再提交 2 个,`2 + 2 > 3 = True` 也会被拒,而如果改成逐个入队的话其实第 1 个是成功的。 这个不是 bug,是策略选择——批量提交要么全过要么全拒,比一半成功一半失败好。没问题。 **真正的竞态问题**:两个并发请求各提交 1 个,用户已有 2 个。预检查都通过(2+1=3 不大于3),然后各自入队,最终用户有 4 个 pending。这个在 DB 层面没有锁。 **建议**: - 短期可以接受,靠兜底的 `safe_enqueue_generation_task` 里再查一次兜底(但也有竞态,只是窗口小一点) - 如果要彻底解决,需要用 DB 行锁或 Redis 原子计数。当前阶段可以先接受"尽力而为"的限流,注释说明一下就行 --- ### ⚠️ P2:预检查逻辑重复 **问题**:task_center 里的 `retry_task_by_id` 和 `retry_project_task` 有几乎一模一样的预检查代码(用户级+全局级),generation_tasks.py 里也有类似的预检查。总共 3 处重复。 **建议**: - 把预检查抽成一个公共函数,比如 `enforce_queue_limits(user_id, repo, extra_count=1)`,超限直接抛 HTTPException - 或者把限流异常统一在中间件/异常处理器里转成 HTTP 响应,路由层就不用 try/except 了 --- ### 总结 | 级别 | 问题 | 影响 | |------|------|------| | P1 | 限流阈值硬编码 + 判断条件不统一 | 维护成本高,边界行为不一致 | | P1 | 并发竞态条件 | 极端情况下可能突破限流 | | P2 | 预检查逻辑 3 处重复 | 维护成本高 | P1 问题修复后即可合并。P2 是代码质量优化,可以后续重构。
xiaoxia added 1 commit 2026-07-11 16:25:06 +08:00
- 阈值统一管理:USER_PENDING_LIMIT/GLOBAL_PENDING_LIMIT 抽到 task_enqueue.py 常量,所有入口引用
- 边界判断统一:预检查用 >=(任务创建前),入队检查用 >(包含当前任务),语义一致
- 并发竞态兜底:发送Celery后再查一次DB计数,超限则回滚任务为failed
- 新增 12 个单元测试(入队后兜底5个 + 边界验证7个)
Author
Owner

复审结论 通过,可以合并

P1 问题全部修复到位,逐一验证如下:

P1-1 阈值统一管理 + 边界判断统一

  • 阈值已抽到 task_enqueue.py 顶部常量:USER_PENDING_LIMIT = 3GLOBAL_PENDING_LIMIT = 20
  • 所有业务入口(generation_tasks 批量/重试、task_center 两级重试、edit_plans)全部引用常量,无硬编码
  • 边界语义统一且自洽:
    • 预检查(任务未创建):>= limit 拒绝
    • 入队前检查(任务已在 DB pending,计数包含自身):> limit 拒绝
    • 两者等价,都是「达到上限就拒绝新任务」
    • 注释写清楚了调用时机和判断条件的对应关系,不会误读

P1-2 并发竞态兜底

  • 采用 DB 层兜底方案:send_task 成功后再查一次计数
  • 超限后调用 _mark_task_failed_safely() 回滚状态为 failed 并抛出异常
  • 状态更新失败有 try/except 兜底,只打日志不崩溃
  • 单测 TestPostEnqueueFinalCheck 覆盖了全局/用户超限回滚、优先级、计数不变等场景

单测覆盖

350+ 行单测,分三组:

  • TestCheckQueueLimits:预检查 8 个用例(正常/用户超限/达到上限/全局超限/优先级/空user_id等)
  • TestSafeEnqueueWithLimits:安全入队 10 个用例(正常/被拒+标记failed/边界值/无user_id/更新失败不崩溃等)
  • TestPostEnqueueFinalCheck:入队后兜底 4 个用例(并发模拟)

P2 遗留(不阻塞合并,后续可单独重构)

  1. 预检查逻辑 4 处重复,可统一封装成 check_queue_limits_or_429() 之类的 helper
  2. edit_plansrender_edit_plan 任务只有预检查,不走 safe_enqueue(因为是另一类任务,不是 generation_task),如果后续也要兜底需要单独处理

整体质量不错,P1 都修干净了,可以 squash merge。

## 复审结论 ✅ 通过,可以合并 P1 问题全部修复到位,逐一验证如下: ### P1-1 阈值统一管理 + 边界判断统一 ✅ - 阈值已抽到 `task_enqueue.py` 顶部常量:`USER_PENDING_LIMIT = 3`、`GLOBAL_PENDING_LIMIT = 20` - 所有业务入口(generation_tasks 批量/重试、task_center 两级重试、edit_plans)全部引用常量,无硬编码 - 边界语义统一且自洽: - 预检查(任务未创建):`>= limit` 拒绝 - 入队前检查(任务已在 DB pending,计数包含自身):`> limit` 拒绝 - 两者等价,都是「达到上限就拒绝新任务」 - 注释写清楚了调用时机和判断条件的对应关系,不会误读 ### P1-2 并发竞态兜底 ✅ - 采用 DB 层兜底方案:`send_task` 成功后再查一次计数 - 超限后调用 `_mark_task_failed_safely()` 回滚状态为 failed 并抛出异常 - 状态更新失败有 try/except 兜底,只打日志不崩溃 - 单测 `TestPostEnqueueFinalCheck` 覆盖了全局/用户超限回滚、优先级、计数不变等场景 ### 单测覆盖 ✅ 350+ 行单测,分三组: - `TestCheckQueueLimits`:预检查 8 个用例(正常/用户超限/达到上限/全局超限/优先级/空user_id等) - `TestSafeEnqueueWithLimits`:安全入队 10 个用例(正常/被拒+标记failed/边界值/无user_id/更新失败不崩溃等) - `TestPostEnqueueFinalCheck`:入队后兜底 4 个用例(并发模拟) ### P2 遗留(不阻塞合并,后续可单独重构) 1. 预检查逻辑 4 处重复,可统一封装成 `check_queue_limits_or_429()` 之类的 helper 2. `edit_plans` 的 `render_edit_plan` 任务只有预检查,不走 safe_enqueue(因为是另一类任务,不是 generation_task),如果后续也要兜底需要单独处理 --- 整体质量不错,P1 都修干净了,可以 squash merge。
xiaoxia merged commit efe7f6b52a into develop 2026-07-11 16:29:08 +08:00
Sign in to join this conversation.