产品库解析编排 Worker 技术方案¶
- 版本:v1.0
- 日期:2026-05-18
- 作者:Codex
- 适用项目:
Tender_Documents
1. 背景¶
当前产品库已支持上传确认后异步解析,主流程位于:
app/routes/v1/product_library/main.pyapp/services/product_library/service.pyapp/services/arq_queue_service.pyapp/worker/parser_jobs.py
现有逻辑已具备基础后台任务能力,但尚未形成“多阶段编排 + 阶段可恢复 + 图片先理解后上传 OSS”的标准化流水线。对于线上高并发和复杂文件,主要风险是:
- 文本解析、图片提取、图片识别、OSS 上传耦合在单段执行里,可观测性不足。
- 任意阶段失败会导致整次任务失败,缺少明确恢复点。
- 图片处理尚未形成“先分类筛选,再上传 OSS”的一致策略。
- 前端能看到任务状态,但阶段级调试和重试信息不足。
因此,本方案将产品库解析升级为“Orchestrator 编排 + Stage 任务状态 + 结果分层落库”。
2. 目标与非目标¶
2.1 目标¶
- API 仅创建任务并入队,立即返回
parse_job_id。 - Worker 按阶段执行:文本抽取、图片抽取、字段抽取、图片识别、关联、上传、入库。
- 每个阶段有独立状态、耗时、错误和重试记录。
- 图片遵循“先理解再上传 OSS”,只上传有效图片。
- 支持任务取消、失败重试、从失败阶段恢复。
- 保持与当前
ProductParseJob、SSE 事件流兼容,支持灰度演进。
2.2 非目标¶
- 本期不重构前端页面交互样式。
- 本期不做跨任务图片去重和全局素材去重。
- 本期不引入复杂 DAG 调度框架,先基于 ARQ 实现串行阶段编排。
- 本期不改动产品库业务字段定义(仍沿用现有
products表字段)。
3. 一句话结论¶
产品库解析应实现为一个可恢复的后台编排任务:
ProductParseJob (主任务)
-> ProductParseStageRun[] (阶段运行记录)
-> ProductParseImage[] (图片中间产物)
-> Product[] (最终业务结果)
数据库保存状态真相;ARQ 负责调度;Redis Pub/Sub 负责实时事件;OSS 只接收筛选后的有效图片。
4. 总体架构¶
4.1 主流程¶
flowchart LR
A[前端确认上传] --> B[API 创建 ProductParseJob]
B --> C[ARQ 入队 orchestrator]
C --> D[Worker 执行阶段编排]
D --> E1[EXTRACT_TEXT]
D --> E2[EXTRACT_IMAGES]
D --> E3[EXTRACT_FIELDS]
D --> E4[CLASSIFY_IMAGES]
D --> E5[LINK_IMAGE_PRODUCT]
D --> E6[UPLOAD_VALID_IMAGES]
D --> E7[PERSIST_PRODUCTS]
E7 --> F[更新 Job 终态与统计]
F --> G[Redis Pub/Sub 推送事件]
G --> H[前端 SSE 订阅]
4.2 分层职责¶
| 层 | 职责 | 不做什么 |
|---|---|---|
| API | 鉴权、建任务、入队、返回 parse_job_id |
不同步执行全文解析 |
| Orchestrator Worker | 阶段调度、状态推进、失败恢复控制 | 不直接承担复杂模型 Prompt 规则维护 |
| Stage 执行器 | 执行具体阶段逻辑 | 不决定全局编排顺序 |
| DB | 存任务状态、阶段日志、图片中间结果、最终产品 | 不做消息分发 |
| Redis Pub/Sub | 推送实时进度事件 | 不作为权威状态来源 |
| OSS | 存储筛选后的有效图片 | 不存无关临时图片 |
5. 编排阶段设计¶
统一阶段枚举(建议):
INITEXTRACT_TEXTEXTRACT_IMAGESEXTRACT_FIELDSCLASSIFY_IMAGESLINK_IMAGE_PRODUCTUPLOAD_VALID_IMAGESPERSIST_PRODUCTSDONE
阶段输入输出约束:
EXTRACT_TEXT:输入文件;输出source_markdown、表格文本、页数。EXTRACT_IMAGES:输出临时图片清单(路径、页码、尺寸、哈希)。EXTRACT_FIELDS:输出产品字段数组(含confidence)。CLASSIFY_IMAGES:输出图片类型(product_image/product_brochure/product_certificate/irrelevant)和置信度。LINK_IMAGE_PRODUCT:输出“图片 -> 产品键”关联结果。UPLOAD_VALID_IMAGES:仅上传有效图片并回写oss_url。PERSIST_PRODUCTS:事务落库产品、更新任务统计。
6. 数据模型设计¶
6.1 复用 product_parse_jobs(主任务表)¶
当前表已具备基础字段,建议新增或扩展:
stage: 当前阶段(字符串或枚举)progress: 0~100retry_count: 当前重试次数max_retries: 最大重试次数(默认 3)result_summary: JSONB(parsed_count,image_total,image_uploaded,low_confidence_count)model_plan: JSONB(文本模型、视觉模型、兜底模型)finished_at: 结束时间(可复用现有时间字段语义)
状态建议:
PENDING/RUNNING/PARTIAL_SUCCESS/COMPLETED/FAILED/CANCELLED
说明:当前 ProductParseStatus 为 PENDING/PROCESSING/COMPLETED/FAILED/CANCELLED,可先兼容映射,后续平滑升级。
6.2 新增 product_parse_stage_runs¶
用于阶段审计与恢复点记录。
建议字段:
idUUID PKjob_idUUID indexstagestring indexstatusstring (PENDING/RUNNING/SUCCESS/FAILED/SKIPPED/CANCELLED)attemptintstarted_attimestamptzfinished_attimestamptzduration_msintinput_snapshotJSONBoutput_snapshotJSONBerror_codestring nullableerror_messagetext nullable
建议索引:
idx_stage_runs_job_stage_attempt (job_id, stage, attempt desc)idx_stage_runs_job_started (job_id, started_at desc)
6.3 新增 product_parse_images¶
用于图片中间态管理(临时抽取 -> 分类 -> 关联 -> 上传)。
建议字段:
idUUID PKjob_idUUID indexpage_nointtemp_pathtextsha256stringwidthintheightintimage_rolestring (product_image/product_brochure/product_certificate/irrelevant)confidencefloatrelated_product_keystring nullableupload_statusstring (PENDING/SKIPPED/UPLOADED/FAILED)oss_urltext nullablereasontext nullablecreated_at/updated_at
7. 状态机与恢复策略¶
7.1 任务状态机¶
PENDING -> RUNNING -> COMPLETED
PENDING -> RUNNING -> PARTIAL_SUCCESS
PENDING -> RUNNING -> FAILED
PENDING -> RUNNING -> CANCELLED
FAILED -> RUNNING (手动重试)
7.2 阶段状态机¶
PENDING -> RUNNING -> SUCCESS
PENDING -> RUNNING -> FAILED
FAILED -> RUNNING -> SUCCESS (重试)
PENDING/RUNNING -> CANCELLED
7.3 恢复点规则¶
- 每阶段成功后写
output_snapshot。 - 重试时从首个失败阶段恢复,不重跑已成功阶段。
EXTRACT_TEXT/EXTRACT_IMAGES产物可复用,避免重复下载与重复 OCR。UPLOAD_VALID_IMAGES必须幂等:同sha256 + related_product_key不重复上传。
8. 失败重试与错误分层¶
错误分层:
- 不可重试:文件损坏、格式不支持、正文为空。
- 可重试:模型超时、429/5xx、网络中断、OSS 短暂失败。
重试策略(建议):
- 每阶段
max_attempts=3。 - 指数退避:2s / 5s / 15s。
- 超过次数后任务置
FAILED,并记录失败阶段与错误码。
取消策略:
- 用户取消后置
CANCELLED。 - Orchestrator 每阶段开始前检查取消状态。
- 已上传 OSS 的文件不回滚删除;仅保证不继续下游执行。
9. 模型与判定策略¶
建议模型计划:
- 文本字段抽取:
Qwen/Qwen3.6-35B-A3B - 图片分类识别:
Qwen/Qwen3-VL-32B-Instruct - 疑难兜底复核:
Qwen/Qwen3.5-397B-A17B(仅低置信度触发)
低置信度策略(建议):
- 文本字段
confidence < 0.75标记为待复核。 - 图片分类
confidence < 0.70默认不上传 OSS。 image_role=irrelevant一律不上传 OSS。
10. 与现有代码的改造点¶
10.1 路由层¶
文件:app/routes/v1/product_library/main.py
- 保留现有
confirm-upload行为。 - 返回值增加编排信息:
stage,progress。 - 新增“重试任务”接口(可选):
POST /{library_id}/parse/jobs/{parse_job_id}/retry
10.2 服务层¶
文件:app/services/product_library/service.py
- 将
process_parse_file_background拆为 orchestrator + stage executor。 - 抽离阶段函数:
_stage_extract_text_stage_extract_images_stage_extract_fields_stage_classify_images_stage_link_image_product_stage_upload_valid_images_stage_persist_products- 新增阶段日志写入与恢复读取函数。
10.3 队列与 Worker¶
文件:
app/services/arq_queue_service.pyapp/worker/parser_jobs.py
改造建议:
- 新增
enqueue_product_parse_orchestrator(job_id, ...)。 - Worker 入口改为 orchestrator 任务,内部统一推进阶段。
- 预留分布式扩展位:后续可把
CLASSIFY_IMAGES拆成子任务并发执行。
11. 事件流与可观测性¶
SSE 事件建议统一:
status: 主状态变化stage: 当前阶段变化progress: 百分比更新metrics: 阶段统计(如parsed_count,image_uploaded)error: 错误信息done: 终态摘要
日志指标建议:
- 每阶段耗时分布(P50/P95)
- 每阶段失败率
- 模型调用成功率与重试次数
- 单任务图片数与有效上传率
12. 灰度发布方案¶
- 第 1 阶段:保留现有逻辑,新增编排字段与阶段日志表,不改接口契约。
- 第 2 阶段:将任务执行切换到 orchestrator(feature flag 控制)。
- 第 3 阶段:启用“先理解后上传”图片策略。
- 第 4 阶段:启用“从失败阶段恢复重试”。
回滚策略:
- 关闭 feature flag,回退到旧
process_parse_file_background。 - 新增表仅作为旁路日志,不影响主业务读写。
13. 验收标准¶
- 上传后 2 秒内返回
parse_job_id。 - 前端可看到阶段推进与进度变化。
- 文档内无关图片不会进入 OSS 与产品字段。
- 单阶段失败可重试并从失败点恢复。
- 任务取消后不再继续执行后续阶段。
- 任务完成后产品字段与图片字段可直接用于产品库列表展示。
14. 后续迭代¶
- 图片去重:基于
sha256的跨任务复用。 - 图片并发分类:批量子任务化。
- 人工审核回路:低置信度结果进入审核队列。
- 模型 A/B:按租户或文件类型动态路由文本/视觉模型。