跳转至

标书动态解析 Worker 池技术设计

1. 文档目的

本文档记录标书解析队列从“单 Worker 串行执行”优化为“可配置动态 Worker 槽位”后的架构、调度规则、隔离措施、监控接口和运维要点,方便后续排障、扩容和交接。

本文中的“动态”指任务与 Worker 槽位的动态分配,不是容器数量的自动伸缩。Compose 固定声明 4 个 Worker 容器,TENDER_PARSE_WORKER_COUNT 控制其中 1~4 个槽位实际消费任务,默认启用 4 个。

2. 改造背景

改造前,所有用户的标书解析任务由同一个 Worker 串行消费。当某个复杂 .doc 文件在旧转换阶段卡住时,会长时间占用唯一的执行位,导致后续所有用户的任务都无法继续。

主要问题包括:

  • 单个异常文件可以阻塞全局队列。
  • 文档提取子进程可能在任务超时后继续运行,长时间占用一个 CPU 核心。
  • 不同用户之间没有执行隔离。
  • 管理员无法直观查看 Worker 是否在线、忙碌或卡住。
  • ARQ 累计计数与当前任务状态缺少统一的管理端展示。

3. 设计目标

  1. 全局最多允许 4 个标书解析任务并行执行。
  2. 每个 Worker 同一时刻只执行 1 个任务。
  3. 一个用户同一时刻只能有 1 个标书解析任务真正执行。
  4. Worker 槽位不与用户永久绑定,任务完成后立即释放给其他用户。
  5. 异常文件只影响当前槽位,不能卡住整个队列。
  6. 对文档提取子进程设置硬超时,超时或取消时清理整个进程组。
  7. 在管理后台展示 Worker 实时状态和历史解析记录。
  8. 保持现有上传、正文提取、快速解析和 SSE 输出业务契约不变。

4. 总体架构

flowchart LR
    API["API 入队"] --> Q["Redis ARQ 共享队列 tender_parse"]
    Q --> W1["Worker 1\nmax_jobs=1"]
    Q --> W2["Worker 2\nmax_jobs=1"]
    Q --> W3["Worker 3\nmax_jobs=1"]
    Q --> W4["Worker 4\nmax_jobs=1"]

    W1 --> UL["Redis 用户级槽位锁"]
    W2 --> UL
    W3 --> UL
    W4 --> UL

    W1 --> P1["独立文档提取子进程"]
    W2 --> P2["独立文档提取子进程"]
    W3 --> P3["独立文档提取子进程"]
    W4 --> P4["独立文档提取子进程"]

    W1 --> S["Redis Worker 心跳与运行状态"]
    W2 --> S
    W3 --> S
    W4 --> S
    S --> ADMIN["Admin Worker 监控接口"]

启用的 Worker 消费同一个 tender_parse 队列,因此它们是一个共享 Worker 池,而不是多条固定队列。任意槽位空闲时,都可以领取队列中下一个可执行任务。

5. Worker 槽位与任务类型

每个 Worker 配置 max_jobs=1,同一时刻只会执行下列任务中的一个:

任务 职责 主要资源特征
parse_markdown 下载原始文件,提取正文 Markdown 并落库 PDF/DOCX/DOC 解析可能消耗 CPU 和内存
quick_parse_stream 基于正文调用 LLM 生成概览、评分点等结构化结果 主要等待外部 LLM 网络请求

两类任务共享当前启用的槽位。例如,当 Worker 1 执行 parse_markdown 时,它不会同时执行 quick_parse_stream;其他 Worker 仍可并行处理其他用户的正文提取或 LLM 快速解析。

Worker 容器与资源限制:

Worker 容器名/hostname CPU 限制 内存限制
Worker 1 tender-worker-tender-parse 1 CPU 1500 MiB
Worker 2 tender-worker-tender-parse-2 1 CPU 1500 MiB
Worker 3 tender-worker-tender-parse-3 1 CPU 1500 MiB
Worker 4 tender-worker-tender-parse-4 1 CPU 1500 MiB

因此,默认标书解析池最多同时有 4 个重任务在执行,队列中等待的任务不会执行文档提取或 LLM 请求。

6. 动态补位流程

假设 6 个不同用户同时提交任务:

时间 Worker 1 Worker 2 Worker 3 Worker 4 队列
初始 用户 A 用户 B 用户 C 用户 D E、F 等待
A 完成 用户 E 补位 用户 B 用户 C 用户 D F 等待
B 完成 用户 E 用户 F 补位 用户 C 用户 D 空

调度特性:

  • Worker 不会永久归属某个用户。
  • 任务完成、失败或取消后,Worker 槽位立即可以领取新任务。
  • 启用的槽位被占满时,新任务保留在 Redis/ARQ 队列中。
  • ARQ 负责队列消费和延后重试;本实现不额外维护用户到 Worker 的静态映射表。
  • 队列总体接近按入队时间消费,但被用户锁延后的任务不承诺绝对严格 FIFO。

7. 用户级并发限制

所有解析 Worker 共享 Redis 用户锁,以确保同一用户同一时刻只有一个标书解析任务真正执行。

核心 Redis Key:

Key 作用
tender_parse:user_slot:{user_id} 用户级执行锁,通过 SET NX EX 获取
tender_parse:business_tries:{job_id} 只统计真正拿到槽位后的业务执行次数

执行流程:

  1. Worker 领取 ARQ 任务。
  2. 根据 user_id 尝试获取 Redis 用户锁。
  3. 获取成功后,记录真实业务尝试次数并开始执行。
  4. 获取失败表示该用户已有任务在运行,当前任务通过 Retry(defer=ARQ_DEFER) 延后。
  5. 任务在 finally 中使用随机 token 和 Lua 脚本安全释放锁。
  6. 只有锁中 token 与当前任务 token 一致时才删除 Key,避免误删已被其他任务续租或重建的锁。

锁 TTL 根据正文解析和快速解析的最大任务超时计算,并额外增加 60 秒缓冲。即使 Worker 崩溃导致显式释放失败,锁也会在 TTL 到期后自动恢复,不会形成永久死锁。

8. 正文与快速解析的顺序

quick_parse_stream 在执行前会查询正文解析任务状态:

  • 如果 parse_markdown 仍在执行,quick_parse_stream 会延后,不会提前调用 LLM。
  • 正文解析完成后,quick_parse_stream 才能获取用户槽位并执行概览、评分点等提取。
  • 等待正文和等待用户锁的次数不等于业务失败重试次数。
  • 调度等待允许较大的 ARQ max_tries,真正拿到槽位后的业务重试仍由 ARQ_MAX_TRIES 限制。

该设计避免把“正常排队”误判为“业务多次失败”。

9. 文档解析进程隔离

9.1 任务级子进程

parse_markdown 不再把整个文档提取链路直接运行在 ARQ Worker 主进程中,而是启动独立的 document_extract_runner 子进程。

子进程使用 start_new_session=True 创建独立进程组。当任务取消、超时或异常时,外层会:

  1. 向整个进程组发送 SIGTERM。
  2. 最多等待 3 秒正常退出。
  3. 仍未退出时向整个进程组发送 SIGKILL。

因此,即使文档解析过程内部又启动了 Java 派生进程,外层仍能对整条任务进程树进行兜底清理。

9.2 旧版 .doc 的 HWPF 路线

真实 OLE .doc 正文提取使用内嵌 Java HWPF CLI 直接完成。

项目 默认值
Jar 路径 /app/bin/doc-extractor.jar
单次提取硬超时 30 秒
Java 最大堆内存 256 MiB

HWPF 子进程输出段落和表格 block,Python 再统一转换为 Markdown 并估算页数。当 HWPF 超时时,该任务按转换硬超时处理,标记失败并且不重试同一文件,避免反复占用槽位。

9.3 旧版 .xls 的 HSSF 路线

真实 OLE .xls 使用同一个内嵌 Java CLI 的 HSSF 模式直接提取工作表、合并单元格和图片锚点。提取器设置独立临时目录、Java 堆上限和硬超时;超时后终止进程,外层任务进程组再兜底清理。

10. Worker 状态与心跳

每个 Worker 使用容器 hostname 作为默认 Worker ID,因此各容器天然具有独立身份。不应在全局环境中为多个容器设置相同的 ARQ_TENDER_PARSE_WORKER_ID。

关键 Redis Key:

Key 内容
tender_parse:health-check:{worker_id} ARQ 心跳、累计成功/失败/重试/执行数
tender_parse:worker-state:{worker_id} 当前任务 ID、函数名、用户、文件名和开始时间

运行规则:

  • Worker 每 10 秒更新 ARQ 心跳。
  • 任务真正获得用户槽位后才写入忙碌状态。
  • 任务结束时立即清理当前 Worker 状态。
  • 状态删除使用 current_job_id 进行防护,避免旧任务的延迟清理删除新任务状态。
  • Worker 启动和关闭时都会清理已有运行状态。
  • 运行状态 TTL 为最大任务超时加 120 秒,防止崩溃后残留永久的“忙碌”记录。

11. Admin 监控接口

两个接口都需要管理员身份。

11.1 实时 Worker 状态

GET /admin/workers/tender-parse

返回内容包括:

  • 配置、在线、忙碌、空闲和离线 Worker 数。
  • Redis 队列总深度和估算等待数。
  • 每个 Worker 的当前函数、任务、用户、文件名、开始时间和执行时长。
  • ARQ 累计成功、失败、重试和正在执行计数。
  • 心跳 TTL。

queue_depth 包含正在执行和被延后的任务,waiting 是用 queue_depth - busy_count 计算的估算值,不是一个独立的精确队列计数。

Worker 累计计数来自 ARQ 心跳,Worker/容器重启后会重新计数,不应将它当作永久的业务统计。

11.2 历史解析记录

GET /admin/workers/tender-parse/history

支持分页和以下筛选条件:

  • status
  • user_id
  • phone
  • filename
  • task_id
  • started_at
  • ended_at

历史记录来自现有的 IngestionTask、IngestionTaskLog、Tender 和 User 表,不依赖 Worker 心跳,容器重启后仍然保留。本次改造没有新增数据库表或迁移。

12. 配置与默认值

配置 默认值 说明
ARQ_TENDER_PARSE_QUEUE_NAME tender_parse 所有解析 Worker 共享的 ARQ 队列
TENDER_PARSE_WORKER_COUNT 4 实际消费任务的 Worker 槽位数,可通过环境变量设为 1~4
ARQ_TENDER_PARSE_MAX_JOBS 1 单容器同时执行数,生产应保持为 1
ARQ_TENDER_PARSE_WORKER_IDS 内置四个容器名 供 Admin 监控识别预期槽位和离线 Worker
ARQ_TENDER_PARSE_WORKER_ID 空 空时使用 hostname;不要给多个容器全局设置相同值
HWPF_EXTRACTOR_JAR /app/bin/doc-extractor.jar 由 Docker 构建阶段产生并复制进镜像
HWPF_EXTRACT_TIMEOUT_SECONDS 30 OLE .doc 提取硬超时
HWPF_EXTRACT_MAX_HEAP_MB 256 HWPF Java 子进程最大堆内存

TENDER_PARSE_WORKER_COUNT 有代码默认值 4;需要在测试环境降低槽位数时可在 Infisical 或本地 env 覆盖。

13. 功能优点

优点 说明
故障隔离 一个文件占用或失败时,其他槽位仍可服务其他用户
全局并行 不同用户默认最多四路并行,降低高峰期等待时间
用户公平性 同一用户不会同时占用多个槽位
资源可控 单 Worker 单任务且有 CPU/内存限制,可以估算最大并发压力
任务可恢复 用户锁和 Worker 状态都有 TTL,崩溃不会造成永久占用
进程可回收 文档提取通过独立进程组运行,超时时可彻底清理派生进程
可观测 后台可查看槽位、当前任务、心跳、队列与历史记录
兼容老流程 入队、任务表、SSE 和快速解析的外部业务契约保持不变

14. 上线验收

验收项 预期结果
容器状态 4 个 worker-tender-parse* 容器全部 healthy
Admin 实时状态 空闲时 online=4、idle=4、busy=0
四用户并发 4 个 Worker 分别执行 1 个任务
六用户并发 4 个执行、2 个等待,槽位释放后自动补位
同用户并发 只有 1 个任务执行,其余任务延后
PDF/DOCX 回归 原有正文和快速解析流程正常
OLE .doc 走 HWPF 路线完成提取,不先转为 DOCX
异常 .doc 按硬超时失败并释放 Worker,不循环重试
进程清理 任务结束后无长时间满核的残留 Java 进程
历史接口 分页、状态、用户、文件名和时间筛选正常

15. 运维与排障

建议优先通过 Admin 页面查看槽位和队列,再使用以下命令定位容器内部进程:

docker stats --no-stream tender-worker-tender-parse

docker top tender-worker-tender-parse \
  -eo pid,ppid,pcpu,pmem,etime,stat,args

如果某个 Worker 长时间忙碌:

  1. 在 Admin 页查看当前函数、用户、任务 ID、文件名和执行时长。
  2. 通过 docker top 确认是 Python 还是 Java 在消耗 CPU。
  3. 检查任务是否已超过对应的业务超时。
  4. 检查 Worker 日志中是否有进程组清理或用户锁释放异常。
  5. 单个 Worker 重启不会阻塞其他槽位;其用户锁和状态也会通过 TTL 兜底恢复。

16. 已知边界

  • 当前固定声明 4 个容器,可通过 TENDER_PARSE_WORKER_COUNT=1..4 启停消费槽位,但不会根据队列长度自动扩缩容。
  • 所有 Worker 共享一条 ARQ 队列,不是每个用户一条物理队列。
  • 用户锁保证的是“同用户最多 1 个任务执行”,不保证多用户之间绝对严格 FIFO。
  • Worker 累计计数会在容器重启后重置;永久历史应以历史解析接口和数据库为准。
  • LLM 请求期间 Worker 槽位仍被占用,即使此时 CPU 使用率不高。这是为了保持单任务边界和降低并发放大风险。
  • 4 核 8 GiB 机器上应继续保持单 Worker max_jobs=1,扩大 Worker 数量前需要重新评估 CPU、内存和异常文件样本。

17. 相关代码

文件 职责
app/worker/tender_parse_jobs.py 标书解析 Worker 定义、心跳、任务注册和并发数
app/worker/parser_jobs.py parse_markdown 和 quick_parse_stream 执行流程
app/services/tender/parse_user_slot.py 用户级 Redis 锁、延后调度和业务重试计数
app/services/tender/parse_worker_monitor.py Worker 心跳、运行状态和队列快照
app/services/tender/document_extract_process.py 文档提取子进程及进程组清理
app/services/parse/document/parsers/legacy_doc_parser.py OLE .doc 的 HWPF 提取和硬超时
app/services/parse/document/parsers/legacy_xls_parser.py OLE .xls 的 HSSF 提取和硬超时
app/routes/admin/workers.py Worker 实时状态和历史解析接口
docker-compose.app.yaml 4 个 Worker 容器、环境开关和资源限制