跳转至

产品库解析编排 Worker 技术方案

  • 版本:v1.0
  • 日期:2026-05-18
  • 作者:Codex
  • 适用项目:Tender_Documents

1. 背景

当前产品库已支持上传确认后异步解析,主流程位于:

  1. app/routes/v1/product_library/main.py
  2. app/services/product_library/service.py
  3. app/services/arq_queue_service.py
  4. app/worker/parser_jobs.py

现有逻辑已具备基础后台任务能力,但尚未形成“多阶段编排 + 阶段可恢复 + 图片先理解后上传 OSS”的标准化流水线。对于线上高并发和复杂文件,主要风险是:

  1. 文本解析、图片提取、图片识别、OSS 上传耦合在单段执行里,可观测性不足。
  2. 任意阶段失败会导致整次任务失败,缺少明确恢复点。
  3. 图片处理尚未形成“先分类筛选,再上传 OSS”的一致策略。
  4. 前端能看到任务状态,但阶段级调试和重试信息不足。

因此,本方案将产品库解析升级为“Orchestrator 编排 + Stage 任务状态 + 结果分层落库”。


2. 目标与非目标

2.1 目标

  1. API 仅创建任务并入队,立即返回 parse_job_id
  2. Worker 按阶段执行:文本抽取、图片抽取、字段抽取、图片识别、关联、上传、入库。
  3. 每个阶段有独立状态、耗时、错误和重试记录。
  4. 图片遵循“先理解再上传 OSS”,只上传有效图片。
  5. 支持任务取消、失败重试、从失败阶段恢复。
  6. 保持与当前 ProductParseJob、SSE 事件流兼容,支持灰度演进。

2.2 非目标

  1. 本期不重构前端页面交互样式。
  2. 本期不做跨任务图片去重和全局素材去重。
  3. 本期不引入复杂 DAG 调度框架,先基于 ARQ 实现串行阶段编排。
  4. 本期不改动产品库业务字段定义(仍沿用现有 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. 编排阶段设计

统一阶段枚举(建议):

  1. INIT
  2. EXTRACT_TEXT
  3. EXTRACT_IMAGES
  4. EXTRACT_FIELDS
  5. CLASSIFY_IMAGES
  6. LINK_IMAGE_PRODUCT
  7. UPLOAD_VALID_IMAGES
  8. PERSIST_PRODUCTS
  9. DONE

阶段输入输出约束:

  1. EXTRACT_TEXT:输入文件;输出 source_markdown、表格文本、页数。
  2. EXTRACT_IMAGES:输出临时图片清单(路径、页码、尺寸、哈希)。
  3. EXTRACT_FIELDS:输出产品字段数组(含 confidence)。
  4. CLASSIFY_IMAGES:输出图片类型(product_image/product_brochure/product_certificate/irrelevant)和置信度。
  5. LINK_IMAGE_PRODUCT:输出“图片 -> 产品键”关联结果。
  6. UPLOAD_VALID_IMAGES:仅上传有效图片并回写 oss_url
  7. PERSIST_PRODUCTS:事务落库产品、更新任务统计。

6. 数据模型设计

6.1 复用 product_parse_jobs(主任务表)

当前表已具备基础字段,建议新增或扩展:

  1. stage: 当前阶段(字符串或枚举)
  2. progress: 0~100
  3. retry_count: 当前重试次数
  4. max_retries: 最大重试次数(默认 3)
  5. result_summary: JSONB(parsed_count, image_total, image_uploaded, low_confidence_count
  6. model_plan: JSONB(文本模型、视觉模型、兜底模型)
  7. finished_at: 结束时间(可复用现有时间字段语义)

状态建议:

PENDING/RUNNING/PARTIAL_SUCCESS/COMPLETED/FAILED/CANCELLED

说明:当前 ProductParseStatusPENDING/PROCESSING/COMPLETED/FAILED/CANCELLED,可先兼容映射,后续平滑升级。

6.2 新增 product_parse_stage_runs

用于阶段审计与恢复点记录。

建议字段:

  1. id UUID PK
  2. job_id UUID index
  3. stage string index
  4. status string (PENDING/RUNNING/SUCCESS/FAILED/SKIPPED/CANCELLED)
  5. attempt int
  6. started_at timestamptz
  7. finished_at timestamptz
  8. duration_ms int
  9. input_snapshot JSONB
  10. output_snapshot JSONB
  11. error_code string nullable
  12. error_message text nullable

建议索引:

  1. idx_stage_runs_job_stage_attempt (job_id, stage, attempt desc)
  2. idx_stage_runs_job_started (job_id, started_at desc)

6.3 新增 product_parse_images

用于图片中间态管理(临时抽取 -> 分类 -> 关联 -> 上传)。

建议字段:

  1. id UUID PK
  2. job_id UUID index
  3. page_no int
  4. temp_path text
  5. sha256 string
  6. width int
  7. height int
  8. image_role string (product_image/product_brochure/product_certificate/irrelevant)
  9. confidence float
  10. related_product_key string nullable
  11. upload_status string (PENDING/SKIPPED/UPLOADED/FAILED)
  12. oss_url text nullable
  13. reason text nullable
  14. created_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 恢复点规则

  1. 每阶段成功后写 output_snapshot
  2. 重试时从首个失败阶段恢复,不重跑已成功阶段。
  3. EXTRACT_TEXT/EXTRACT_IMAGES 产物可复用,避免重复下载与重复 OCR。
  4. UPLOAD_VALID_IMAGES 必须幂等:同 sha256 + related_product_key 不重复上传。

8. 失败重试与错误分层

错误分层:

  1. 不可重试:文件损坏、格式不支持、正文为空。
  2. 可重试:模型超时、429/5xx、网络中断、OSS 短暂失败。

重试策略(建议):

  1. 每阶段 max_attempts=3
  2. 指数退避:2s / 5s / 15s。
  3. 超过次数后任务置 FAILED,并记录失败阶段与错误码。

取消策略:

  1. 用户取消后置 CANCELLED
  2. Orchestrator 每阶段开始前检查取消状态。
  3. 已上传 OSS 的文件不回滚删除;仅保证不继续下游执行。

9. 模型与判定策略

建议模型计划:

  1. 文本字段抽取:Qwen/Qwen3.6-35B-A3B
  2. 图片分类识别:Qwen/Qwen3-VL-32B-Instruct
  3. 疑难兜底复核:Qwen/Qwen3.5-397B-A17B(仅低置信度触发)

低置信度策略(建议):

  1. 文本字段 confidence < 0.75 标记为待复核。
  2. 图片分类 confidence < 0.70 默认不上传 OSS。
  3. image_role=irrelevant 一律不上传 OSS。

10. 与现有代码的改造点

10.1 路由层

文件:app/routes/v1/product_library/main.py

  1. 保留现有 confirm-upload 行为。
  2. 返回值增加编排信息:stage, progress
  3. 新增“重试任务”接口(可选):
  4. POST /{library_id}/parse/jobs/{parse_job_id}/retry

10.2 服务层

文件:app/services/product_library/service.py

  1. process_parse_file_background 拆为 orchestrator + stage executor。
  2. 抽离阶段函数:
  3. _stage_extract_text
  4. _stage_extract_images
  5. _stage_extract_fields
  6. _stage_classify_images
  7. _stage_link_image_product
  8. _stage_upload_valid_images
  9. _stage_persist_products
  10. 新增阶段日志写入与恢复读取函数。

10.3 队列与 Worker

文件:

  1. app/services/arq_queue_service.py
  2. app/worker/parser_jobs.py

改造建议:

  1. 新增 enqueue_product_parse_orchestrator(job_id, ...)
  2. Worker 入口改为 orchestrator 任务,内部统一推进阶段。
  3. 预留分布式扩展位:后续可把 CLASSIFY_IMAGES 拆成子任务并发执行。

11. 事件流与可观测性

SSE 事件建议统一:

  1. status: 主状态变化
  2. stage: 当前阶段变化
  3. progress: 百分比更新
  4. metrics: 阶段统计(如 parsed_count, image_uploaded
  5. error: 错误信息
  6. done: 终态摘要

日志指标建议:

  1. 每阶段耗时分布(P50/P95)
  2. 每阶段失败率
  3. 模型调用成功率与重试次数
  4. 单任务图片数与有效上传率

12. 灰度发布方案

  1. 第 1 阶段:保留现有逻辑,新增编排字段与阶段日志表,不改接口契约。
  2. 第 2 阶段:将任务执行切换到 orchestrator(feature flag 控制)。
  3. 第 3 阶段:启用“先理解后上传”图片策略。
  4. 第 4 阶段:启用“从失败阶段恢复重试”。

回滚策略:

  1. 关闭 feature flag,回退到旧 process_parse_file_background
  2. 新增表仅作为旁路日志,不影响主业务读写。

13. 验收标准

  1. 上传后 2 秒内返回 parse_job_id
  2. 前端可看到阶段推进与进度变化。
  3. 文档内无关图片不会进入 OSS 与产品字段。
  4. 单阶段失败可重试并从失败点恢复。
  5. 任务取消后不再继续执行后续阶段。
  6. 任务完成后产品字段与图片字段可直接用于产品库列表展示。

14. 后续迭代

  1. 图片去重:基于 sha256 的跨任务复用。
  2. 图片并发分类:批量子任务化。
  3. 人工审核回路:低置信度结果进入审核队列。
  4. 模型 A/B:按租户或文件类型动态路由文本/视觉模型。