标书动态解析 Worker 池技术设计¶
1. 文档目的¶
本文档记录标书解析队列从“单 Worker 串行执行”优化为“可配置动态 Worker 槽位”后的架构、调度规则、隔离措施、监控接口和运维要点,方便后续排障、扩容和交接。
本文中的“动态”指任务与 Worker 槽位的动态分配,不是容器数量的自动伸缩。Compose 固定声明 4 个 Worker 容器,TENDER_PARSE_WORKER_COUNT 控制其中 1~4 个槽位实际消费任务,默认启用 4 个。
2. 改造背景¶
改造前,所有用户的标书解析任务由同一个 Worker 串行消费。当某个复杂 .doc 文件在旧转换阶段卡住时,会长时间占用唯一的执行位,导致后续所有用户的任务都无法继续。
主要问题包括:
- 单个异常文件可以阻塞全局队列。
- 文档提取子进程可能在任务超时后继续运行,长时间占用一个 CPU 核心。
- 不同用户之间没有执行隔离。
- 管理员无法直观查看 Worker 是否在线、忙碌或卡住。
- ARQ 累计计数与当前任务状态缺少统一的管理端展示。
3. 设计目标¶
- 全局最多允许 4 个标书解析任务并行执行。
- 每个 Worker 同一时刻只执行 1 个任务。
- 一个用户同一时刻只能有 1 个标书解析任务真正执行。
- Worker 槽位不与用户永久绑定,任务完成后立即释放给其他用户。
- 异常文件只影响当前槽位,不能卡住整个队列。
- 对文档提取子进程设置硬超时,超时或取消时清理整个进程组。
- 在管理后台展示 Worker 实时状态和历史解析记录。
- 保持现有上传、正文提取、快速解析和 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} |
只统计真正拿到槽位后的业务执行次数 |
执行流程:
- Worker 领取 ARQ 任务。
- 根据
user_id尝试获取 Redis 用户锁。 - 获取成功后,记录真实业务尝试次数并开始执行。
- 获取失败表示该用户已有任务在运行,当前任务通过
Retry(defer=ARQ_DEFER)延后。 - 任务在
finally中使用随机 token 和 Lua 脚本安全释放锁。 - 只有锁中 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 创建独立进程组。当任务取消、超时或异常时,外层会:
- 向整个进程组发送
SIGTERM。 - 最多等待 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 状态¶
返回内容包括:
- 配置、在线、忙碌、空闲和离线 Worker 数。
- Redis 队列总深度和估算等待数。
- 每个 Worker 的当前函数、任务、用户、文件名、开始时间和执行时长。
- ARQ 累计成功、失败、重试和正在执行计数。
- 心跳 TTL。
queue_depth 包含正在执行和被延后的任务,waiting 是用 queue_depth - busy_count 计算的估算值,不是一个独立的精确队列计数。
Worker 累计计数来自 ARQ 心跳,Worker/容器重启后会重新计数,不应将它当作永久的业务统计。
11.2 历史解析记录¶
支持分页和以下筛选条件:
statususer_idphonefilenametask_idstarted_atended_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 长时间忙碌:
- 在 Admin 页查看当前函数、用户、任务 ID、文件名和执行时长。
- 通过
docker top确认是 Python 还是 Java 在消耗 CPU。 - 检查任务是否已超过对应的业务超时。
- 检查 Worker 日志中是否有进程组清理或用户锁释放异常。
- 单个 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 容器、环境开关和资源限制 |