# SupplyAgent 端到端产品需求文档
> 文档性质:现状还原 + 目标产品要求 + 风险整改清单
> 版本:v1.0
> 更新日期:2026-07-24
> 代码基线:`dev-find` / `5ea5e65`
> 核心入口:`supply_infra/scheduler/jobs/run_supply_pipeline.py`
---
## 1. 文档目的
本文用于统一说明 SupplyAgent 当前项目的完整业务流程,重点回答:
1. 每日定时任务如何贯穿 ODPS、MySQL、业务 Agent、找片和 AIGC;
2. 每一阶段的输入、处理规则、输出、状态与失败方式;
3. API、前端和人工补偿入口如何消费流水线结果;
4. 当前实现中已经发现的产品问题、数据风险、工程风险和安全风险;
5. 下一阶段应达到的产品要求、优先级和验收标准。
本文以当前代码实际行为为准。仓库中的 `README.md`、`ARCHITECTURE.md`、`prd/` 和
`zhangbo.md` 部分内容描述的是旧流程或目标蓝图,不能代替本文对现状的说明。
本次分析未触发 ODPS、外部搜索、AIGC 发布等线上副作用;运行态数据量、接口成功率和真实耗时仍需结合生产监控补充。
---
## 2. 产品概述
### 2.1 产品定位
SupplyAgent 是一套面向内容供给的每日需求处理系统。它将多来源需求信号组织到全局分类树中,
结合先验热度和真实效果完成需求分级,再从已有视频点位中扩展搜索意图,自动寻找适合目标受众的
短视频,并把保留候选分发到 AIGC 生产计划。
当前产品形态由五部分组成:
- 通用 Agent 框架:LLM、Tool Calling、Skills、运行日志;
- 数据基础设施:ODPS、MySQL、OSS;
- 每日供给流水线:同步、归类、聚合、分级、拓展、找片、分发;
- FastAPI 查询服务;
- Vue 需求地图和需求/视频证据页面。
### 2.2 核心业务目标
- 每天形成一份可追溯的需求数据快照;
- 将需求词稳定挂靠到全局分类树;
- 基于来源热度和真实效果形成 S/A/B/C/D 需求优先级;
- 为高优需求补充可执行的搜索点位;
- 自动发现与需求相关、偏老年受众且具有分享价值的视频;
- 将合格视频安全、准确、幂等地交给正确的 AIGC 生产链路;
- 通过 API、前端、日志和执行记录解释每个结果是如何产生的。
### 2.3 当前非目标
当前代码尚未完整实现以下目标:
- 需求的多分类挂靠和正式语义关系图;
- 统一的“平台需求”版本化产品对象;
- 人工反馈和线上效果自动回流到次日决策;
- 真正调用发布计划完成内容发布;
- 在页面中展示 `find_agent` 最终候选及完整审计过程;
- 跨实例的任务编排、分布式锁和可靠消息机制。
---
## 3. 用户与外部系统
| 角色/系统 | 主要诉求 | 当前交互 |
|---|---|---|
| 内容策略/运营 | 看懂需求热度、等级、原因和内容证据 | Vue 需求地图、需求列表、视频点位 |
| 数据/算法人员 | 核对数据口径、分类、热度和后验 | MySQL、ODPS、Agent 日志 |
| 运维/开发人员 | 观察任务是否运行、定位失败、手动补偿 | Scheduler 状态接口、日志、`python -m supply_infra.scheduler.jobs.*` |
| ODPS | 提供分类树、需求池、视频解析、ROV/VOV | 定时查询 |
| MySQL | 保存业务快照、关系、执行状态和候选 | 全流程状态底座 |
| OpenRouter/模型服务 | 执行归类、分级、拓展和找片判断 | Agent Tool Calling |
| 抖音搜索/画像/TikHub/千问 | 提供找片、详情、画像和内容解析证据 | `find_agent` 外部工具 |
| OSS | 保存 Agent `.log`、`.jsonl` 和 HTML 可视化 | 每次 Agent 运行后上传 |
| AIGC 平台 | 创建视频爬取计划并绑定生产计划 | 流水线末端调用 |
---
## 4. 关键业务对象
| 对象 | 说明 | 核心数据表 |
|---|---|---|
| 全局分类树 | 稳定的内容分类骨架 | `global_tree_category` |
| 树下元素 | ODPS 提取的实质元素及挂靠 | `global_tree_element` |
| 策略需求池 | 每日多策略需求原始行 | `multi_demand_pool_di` |
| 需求归类词 | 由需求名称拆出的词及单一挂靠节点 | `demand_belong_category` |
| 需求池匹配边 | 归类词与需求池行的子串关系 | `demand_belong_pool_rel` |
| 词级热度 | 四项先验和两项后验的 avg/count | `demand_popularity_stats` |
| 树节点热度 | 六维向祖先汇总后的节点指标 | `category_tree_weight` |
| 分级计划 | 分类节点分组及任务明细 | `demand_grade_plan*` |
| 需求等级 | 每日 S/A/B/C/D 结果 | `demand_grade`、`demand_grade_category_rel` |
| 源视频及点位 | 需求池视频的标题、选题、三类点位 | `multi_demand_video_detail`、`multi_demand_video_point` |
| 需求拓展 | S/A 需求从视频点位选出的拓展意图 | `demand_video_expansion*` |
| 找片运行 | 搜索树、候选证据、评分和分池 | `video_discovery_run/search/candidate` |
| AIGC 分发状态 | 候选被分配到的爬取/生产/发布计划标识 | `video_discovery_candidate` |
| 任务执行记录 | 总流水线 started/finished/failed/skipped | `scheduler_job_execution` |
| Agent 审计日志 | 模型输入、输出、工具调用及 HTML | 本地 `logs/`、OSS、`oss_logs` |
---
## 5. 每日总流程
```mermaid
flowchart LR
START["API/CLI 启动 Scheduler"] --> CRON["每日 15:00
Asia/Shanghai"]
CRON --> T1["① 同步 T-1 全局分类树"]
T1 --> POOL["② 同步当日策略需求池"]
POOL --> CLASSIFY["归类新词并建立匹配边"]
CLASSIFY --> METRIC["回填 ROV/VOV
计算词级和树级热度"]
METRIC --> SOURCE_VIDEO["同步源视频标题与三类点位"]
SOURCE_VIDEO --> GRADE["③ 需求分级
S/A/B/C/D"]
GRADE --> EXPAND["④ S/A 视频点位拓展"]
EXPAND --> FIND["⑤ Top 200 需求找片"]
FIND --> AIGC["⑥ 创建 AIGC 爬取计划
绑定生产计划"]
AIGC --> API["FastAPI 查询"]
API --> WEB["Vue 需求地图/需求/视频证据"]
```
### 5.1 业务日期
- 未显式传入 `biz_dt` 时,按 `SCHEDULER_TIMEZONE` 的当天生成 `YYYYMMDD`;
- 全局分类树使用 `biz_dt - 1 天` 的 ODPS 分区;
- 策略需求池、热度、分级、拓展、找片使用 `biz_dt`;
- 视频解析同步没有使用传入的 `biz_dt`,而是固定读取任务实际执行日的昨天。
---
## 6. 调度与运行控制
### 6.1 启动方式
1. FastAPI 启动时执行 `init_db()`;
2. 若 `SCHEDULER_ENABLED=true`,在 API lifespan 中启动后台 Scheduler(唯一启动入口);
3. 可通过 `python -m supply_infra.scheduler.jobs.run_supply_pipeline [YYYYMMDD]` 手动执行全链路。
### 6.2 调度参数
| 参数 | 当前值 |
|---|---|
| 触发时间 | 每日 15:00 |
| 时区 | 默认 `Asia/Shanghai` |
| Job ID | `run_supply_pipeline` |
| 同一 Scheduler 最大实例 | 1 |
| 合并错过的执行 | `coalesce=True` |
| 允许延迟 | 3600 秒 |
| 找片需求上限 | 200 |
| 找片并发 | 2 |
| 分级并发 | 5 |
| 点位拓展并发 | 5 |
### 6.3 当前异常策略
- 总流水线用进程内 `threading.Lock` 防止同进程重入;
- 每个主步骤捕获异常后记录失败,但继续执行后续步骤;
- 所有步骤完成后,只有全部步骤成功才把总运行标记为成功;
- 总运行的开始和结束分别写入 `scheduler_job_execution`;
- 子步骤没有独立、统一的持久化执行记录;
- Agent 运行另有本地日志、JSONL、HTML 和 OSS 记录。
---
## 7. 阶段一:同步全局分类树
### 7.1 输入
- ODPS `pattern_mining_element` 的 T-1 分区;
- ODPS `public_pattern_mining_category` 的 T-1 分区;
- 固定过滤 `execution_id=401`、`source_type/element_type=实质` 等条件。
### 7.2 处理
1. 拉取有效元素及其 `category_id`;
2. 拉取全量分类;
3. 从元素命中的分类向上补齐全部祖先;
4. 按 ODPS `source_id` 去重;
5. 对 MySQL 不存在的分类分配本地 ID;
6. 转换父节点 ID 后 `INSERT IGNORE`;
7. 重载 ID 映射并写入元素。
### 7.3 输出
- `global_tree_category`;
- `global_tree_element`;
- 本阶段的拉取、插入、跳过和失败数量。
### 7.4 当前行为边界
- 已存在分类不会更新名称、描述、层级和父节点;
- ODPS 已删除或迁移的节点不会在 MySQL 自动软删除;
- 已有元素也只追加,不校正旧挂靠;
- 缺失父节点的分类会被跳过,但不会阻断后续流水线。
---
## 8. 阶段二:同步需求池并生成热度底座
本阶段实际上是一个包含七个子阶段的内部流水线。任一子阶段失败时,其余子阶段仍会继续。
### 8.1 子阶段 A:同步当日需求池
输入为 ODPS `dwd_multi_demand_pool_di` 的 `biz_dt` 分区。
处理规则:
1. 先比较 ODPS 与 MySQL 当日 `(strategy, demand_id)` 去重行数;
2. 行数相同时,完全跳过明细拉取;
3. 行数不同时,拉取全量明细并按 `(strategy, demand_id)` 做插入、删除和更新;
4. `video_list` 最多保留 10 个视频;
5. `weight=0` 转为 `NULL`;
6. 对已有行只更新 `video_list`、`video_count` 和 `weight`。
输出为 `multi_demand_pool_di` 当日数据。
### 8.2 子阶段 B:归类新需求词
1. 查询当日全部 `demand_name`;
2. 使用空格拆词,汇总为去重词集合;
3. 过滤 `demand_belong_category` 已存在的名称;
4. 每 100 个词调用一次 `demand_belong_category_agent`;
5. Agent 查询真实分类树,选择分类节点并写入名称、节点和 reason。
当前表以 `name` 唯一,所以一个词实际上只能保存一个分类节点。
### 8.3 子阶段 C:建立归类词与需求池关系
1. 读取全部有效归类词;
2. 读取需求池全表,而非仅当日分区;
3. 若 `归类词.name in 需求池.demand_name`,则建立关系边;
4. 已有边跳过,只插入缺失边;
5. 合并匹配行的视频,去重后最多保留 10 个,回填归类词。
输出:
- `demand_belong_pool_rel`;
- `demand_belong_category.video_list`。
### 8.4 子阶段 D:回填真实 ROV/VOV
1. 查询 `biz_dt-7天` 至 `biz_dt` 的生产效果数据;
2. 仅保留人工 AGC 和自动 AGC;
3. 计算各特征相对全局基线的 `rov_diff`、`vov_diff`;
4. SQL 最多返回 1000 行;
5. 同一特征值有多行时选择 `rov_diff` 优先、再比较 `vov_diff` 的较高记录;
6. 使用 `特征值 == demand_name` 精确匹配回填当日需求池。
输出:
- `multi_demand_pool_di.real_rov_7d`;
- `multi_demand_pool_di.real_vov_7d`。
### 8.5 子阶段 E:计算词级热度
对 `demand_belong_category` 的全部有效词逐一处理:
1. 在当日需求池用 `demand_name LIKE %词%` 查找匹配行;
2. 将来源策略映射为四项先验;
3. 非零权重计算 avg/count;
4. ROV/VOV 按匹配需求名称聚合;
5. 按 `(demand_category_id, biz_dt)` upsert。
策略映射:
| 来源策略 | 指标 |
|---|---|
| 新热事件 | `ext_pop` |
| 逐月 | `plat_sust_pop` |
| 去年同期阳历、去年同期阴历 | `plat_ly_pop` |
| 当下供需gap | `recent_pop` |
输出为 `demand_popularity_stats`。
### 8.6 子阶段 F:计算树节点热度和排名
1. 将同一挂载节点下的词级指标按样本数加权;
2. 每个分类节点汇总其整个子树内所有挂载点;
3. 六维分别保存 avg/count;
4. 四项先验分别做全树排名归一化;
5. 将有效的四维排名分直接相加为 `total_score`;
6. ROV/VOV 保存但不进入 `total_score`。
输出为 `category_tree_weight`。
### 8.7 子阶段 G:同步需求池源视频
1. 从需求池全表所有 `video_list` 收集视频 ID;
2. 已存在于详情表的 ID 直接跳过;
3. 对待处理 ID 分批查询 ODPS 视频解析结果;
4. ODPS 分区固定使用任务执行日的昨天;
5. 提取标题、最终选题、灵感点、目的点和关键点;
6. 写入详情表,并替换本批视频的点位行。
输出:
- `multi_demand_video_detail`;
- `multi_demand_video_point`。
---
## 9. 阶段三:需求分级
### 9.1 目标
将当日需求分为 S/A/B/C/D,供后续点位拓展、找片和资源分配使用。
### 9.2 自动计划
1. 从 `category_tree_weight.hung_word_count > 0` 的节点中找待分配节点;
2. 同分类、同父分类优先组合;
3. 每组目标约 30 条需求;
4. 每日最多 200 个组;
5. 计划及组写入 `demand_grade_plan`、`demand_grade_plan_group`;
6. 通过匹配关系解析当日需求池行,物化为组明细。
当前每日定时任务默认使用代码自动分组,不使用 `demand_grade_orchestrator_agent`。
### 9.3 Agent 分级
1. 5 个 worker 并发领取 pending 组;
2. 每组明细再按最多 30 条拆分;
3. `demand_grade_agent` 查询:
- 需求来源内排名分;
- 分类树 `total_score`;
- 父节点及全部兄弟节点;
- 词级和节点级 ROV/VOV;
- 同名或包含关系的需求池行;
4. Agent 判定 S/A/B/C/D 并写 reason;
5. 保存工具确定性重算 `score`:
- 每个 strategy 内按非零 weight 排名归一化;
- 同一需求已有来源等权平均;
- 结果映射为 0~100;
6. 保存分类关系、关联需求池行、来源、视频列表和先后验快照。
### 9.4 输出
- `demand_grade`;
- `demand_grade_category_rel`;
- 计划组和组明细状态。
---
## 10. 阶段四:S/A 需求视频点位拓展
### 10.1 入选条件
一条需求必须同时满足:
- 当日等级为 S 或 A;
- `demand_grade.video_list` 非空;
- 对应视频已同步出三类点位;
- 当日尚未有 finished 拓展记录。
### 10.2 处理
1. 汇总需求关联视频的全部灵感点、目的点和关键点;
2. 5 个 worker 分别为单条需求调用 `demand_video_expand_agent`;
3. Agent 判断点位是否与原需求存在包含、细分或同意图关系;
4. 过滤与原需求完全相同、过宽、无关或非需求表达的点;
5. 保存拓展文本、点位类型、视频、描述和 reason;
6. 无候选时允许保存 0 条,并将运行记录标记为 finished。
### 10.3 输出
- `demand_video_expansion`;
- `demand_video_expansion_run`。
---
## 11. 阶段五:find_agent 视频发现
### 11.1 入选与排序
1. 只读取 S/A 需求;
2. 只保留已经产生 `demand_video_expansion` 的需求;
3. 使用拓展记录按视频组装参考标题和点位;
4. 当前实现先按 `score` 降序,再用 S/A 作为次级排序;
5. 每天最多取前 200 条;
6. 2 个 worker 并发执行;
7. 当日存在 `running` 或 `finished` 找片记录时默认跳过。
### 11.2 单需求找片流程
1. 系统预创建 `video_discovery_run` 并生成 `run_id`;
2. Agent 根据需求、参考视频和点位形成 2~3 个搜索假设;
3. 调用内部抖音搜索、TikHub 或作者作品搜索;
4. 每个搜索页写入 `video_discovery_search`;
5. 候选视频写入 `video_discovery_candidate`,初始为 `pending_evaluation`;
6. 对高潜候选补充详情、视频点赞用户画像和作者粉丝画像;明确不使用视频理解;
7. Agent 独立判断:
- R:需求相关性;
- E:老年受众倾向;
- S:分享价值;
- V:联合价值;
8. 候选进入:
- `primary`:主推荐;
- `rejected`:淘汰;
- `pending_evaluation`:尚未完成的过程状态,不是最终等级;
9. Agent 保存最终候选评估并将运行置为 `finished`;
10. 数据库审计通过后重新查询最终状态,再输出主推荐、淘汰原因、搜索树和缺失证据;
11. completion guard 校验顺序和报告分池,报告之后不再修改数据库。
### 11.3 输出
- `video_discovery_run`;
- `video_discovery_search`;
- `video_discovery_candidate`;
- Agent 审计日志和 OSS HTML。
---
## 12. 阶段六:AIGC 分发
### 12.1 当前入选规则
- 选择当日 `decision_bucket=primary` 的候选;
- `aweme_id` 必须非空;
- 默认跳过已经写入 `aigc_crawler_plan_id` 的候选;
- 当前查询不要求所属 `video_discovery_run.status=finished`。
### 12.2 当前处理
1. 对 AIGC 计划映射按 `(生成ID, 发布ID)` 去重;
2. 将所有候选轮询均匀分配到所有计划对,不按需求或分类路由;
3. 每 10 个视频创建一个爬取计划;
4. 获取生产计划详情;
5. 将爬取计划追加为生产计划的输入源;
6. 绑定成功后,在候选表保存:
- 爬取计划 ID;
- 生成计划 ID;
- 发布计划 ID;
- 分配标签。
### 12.3 重要语义说明
当前代码只“创建爬取计划并绑定生成计划”,没有使用 `publish_plan_id` 调用发布接口,
也没有验证内容生产或发布完成。因此现阶段准确名称应为“AIGC 生产计划分发”,不能视为真正发布完成。
---
## 13. API 与前端消费流程
### 13.1 API
| 接口 | 当前用途 |
|---|---|
| `GET /health` | 进程健康 |
| `GET /api/scheduler/status` | 当前进程 Scheduler 状态和下次运行时间 |
| `GET /api/category-tree` | 分类树、六维指标、`total_score` |
| `GET /api/demand-belong-category` | 归类词及单一挂靠 |
| `GET /api/demand-belong-category/{id}/videos` | 归类词的源视频和点位 |
| `GET /api/demand-grade` | 当日分级需求 |
| `GET /api/demand-grade/{id}/videos` | 分级需求的拓展视频和点位 |
| `GET /api/video-discovery/demands` | 以需求为中心的分级与拓展摘要 |
| `GET /api/video-discovery/demands/{id}` | 需求的拓展视频证据 |
| `GET /api/demand-belong-oss-logs` | 归类 Agent 日志 |
### 13.2 前端
| 页面 | 当前展示 |
|---|---|
| 平台全局需求地图 | 分类树、热度、分级需求及下钻 |
| 全局分类树 | 可展开分类树和多维热度 |
| 需求归类过程 | 归类 Agent OSS 日志 |
| 需求汇总/视频发现 | 分级需求、源视频、拓展点位 |
当前前端“视频发现”页面没有读取 `video_discovery_candidate`,因此展示的是找片输入证据,
不是 `find_agent` 最终找到的主推荐和淘汰候选。
---
## 14. 旁路与手动流程
以下能力存在于项目中,但不属于每日总流水线:
- `generate_demand_agent`:按单一热度维度写入 `generated_demand`;
- `demand_grade_orchestrator_agent`:模型统筹分组,当前定时路径已改为代码自动分组;
- `python -m supply_infra.scheduler.jobs.grade_demand_pool --retry-failed`:手动重试失败分级明细;
- 通过 API 或 `python -m supply_infra.scheduler.jobs.*` 单独跑各流水线步骤;
- 日志可视化和手动 OSS 上传。
`generated_demand` 当前没有接入总流水线、API 或前端,属于孤立产物。
---
## 15. 功能需求
### FR-01 每日批次
- 系统必须为每个 `biz_dt` 建立唯一每日批次;
- 必须记录使用的分区、代码版本、配置版本、模型、Prompt 版本和每步状态;
- 同一业务日重跑必须生成明确的重跑版本或复用幂等键,不得静默覆盖。
### FR-02 依赖门禁
- 下游步骤只能消费通过完整性校验的上游产物;
- 分类树或需求池失败时,不得执行依赖其结果的分级;
- 分级覆盖不完整时,不得将该日结果作为可发布批次;
- 找片未完成审计时,不得进入 AIGC 分发。
### FR-03 数据同步
- 同步判断必须基于主键集合和内容校验值,不得只比较行数;
- 分类、元素、需求池和关系都必须明确支持新增、更新、删除/失效;
- 所有时间窗口和分区必须由 `biz_dt` 派生,支持历史重跑;
- 缺失、0、NULL 和迟到数据必须有不同语义。
### FR-04 需求模型
- 原始需求行、标准需求词和平台需求必须分层;
- 一个需求允许挂靠多个分类节点;
- 每条挂靠和关系必须保存来源、reason、置信度和版本;
- 字符串包含只能作为候选匹配,不得直接成为权威语义关系。
### FR-05 分级
- 每条当日有效需求必须得到“成功、失败、跳过及原因”之一;
- 计划覆盖率和结果覆盖率必须分别达到 100% 才能标记完成;
- S/A/B/C/D 和数值 score 的口径必须版本化;
- Agent 返回成功后必须查询数据库验证全部输入项均已落库;
- 失败项自动重试达到上限后进入补偿队列和报警。
### FR-06 点位拓展
- 每条 S/A 需求必须有明确结果:有拓展、无拓展、无源视频、无点位或执行失败;
- 0 条拓展必须是经工具确认的有效业务结果,不能由“未调用保存工具”推断;
- 没有拓展点的 S/A 需求仍应允许使用原需求进入找片。
### FR-07 视频发现
- 入选顺序必须符合确认后的业务策略,S/A 优先级与 score 排序不得冲突;
- 每个运行必须有租约和超时,崩溃遗留的 `running` 可自动恢复;
- 完成前必须校验搜索页、候选评估、证据、审计和最终状态;
- 强制重跑必须新建 attempt 或清理旧子记录,不能混用两次搜索状态;
- 相同视频跨需求发现时必须全局去重或形成“一视频多需求”关系。
### FR-08 AIGC 分发与发布
- 只有 finished 且审计通过的找片运行可以分发;
- 仅 `primary` 允许自动分发;
- 候选必须按需求分类或明确的路由规则进入正确计划;
- 创建、绑定、生产、发布每个外部动作必须有幂等键和独立状态;
- 必须区分“已分发、已生产、已发布、发布失败”;
- 只有真实发布接口成功并经查询确认后,才能标记“已发布”。
### FR-09 查询与运营
- 前端必须展示每日批次及各阶段状态;
- 必须展示 `find_agent` 的主推荐、淘汰候选和证据;
- 支持按业务日、需求、分类、运行状态和发布状态检索;
- 支持安全的失败项重试,并保留操作审计。
---
## 16. 非功能需求
### 16.1 可靠性
- 调度采用数据库锁、Redis 锁或独立 Scheduler Leader,保证跨进程/副本单实例执行;
- 数据写入和外部调用采用幂等设计;
- 每个步骤支持断点续跑;
- 任一步失败不得造成错误数据进入不可逆的下游系统。
### 16.2 可观测性
- 每步记录开始、结束、输入量、输出量、失败量、耗时和重试次数;
- 提供批次级和需求级状态查询;
- P0/P1 失败应主动报警;
- 日志不得包含 API Token、密码、Cookie 或完整敏感请求体。
### 16.3 性能
- 避免对全部归类词和需求池做 Python 双重循环;
- 高频查询建立明确唯一键和索引;
- API 大结果支持分页、缓存或按子树读取;
- LLM 调用并发必须受全局限流、超时和成本预算控制。
### 16.4 安全
- 生产 API 必须有认证、授权和访问审计;
- 外部 Token 只能通过密钥管理系统注入;
- OSS 日志访问应按敏感级别控制;
- 所有日志和错误信息必须脱敏。
---
## 17. 已发现问题与风险
### 17.1 P0:上线前必须处理
| ID | 问题 | 影响 | 代码现状/证据 | 产品要求 |
|---|---|---|---|---|
| P0-01 | 上游失败后仍继续下游 | 可能用旧树、空需求池或不完整分级继续找片和分发 | 总流水线逐步捕获异常并无条件继续 | 建立依赖 DAG 和硬门禁 |
| P0-02 | 分级存在失败组时仍可能判定完成 | 部分需求未分级,但总步骤返回成功 | `get_execution_snapshot` 的 `execution_complete` 只统计 pending/running,不统计 failed | failed 必须阻止完成 |
| P0-03 | 无计划或无分级也可能成功 | 当日需求完全未处理仍进入拓展和找片 | 自动计划异常被吞掉;空计划快照可被视为 complete | 校验计划覆盖率和分级覆盖率 |
| P0-04 | Agent 返回即把分级明细标记 finished | 模型未保存、少保存或保存工具报错时产生假完成 | worker 不核对 `demand_grade` 实际落库覆盖 | 每批结束后按输入逐条验库 |
| P0-05 | 点位拓展把“未保存”误判为“零结果” | S/A 需求被永久标记完成并从找片链路消失 | 从工具文本解析不到数量时默认为 0,仍写 finished | 必须验证保存工具调用和 run 状态 |
| P0-06 | 找片确定性完成控制(已修复) | 防止在搜索、证据、评估或审计未完成时提前结束 | 已注册数据库审计并启用 completion guard,强制审计后查询最终状态再报告 | 保持顺序与负向回归测试 |
| P0-07 | AIGC 分发不校验找片运行状态 | running/failed 运行中的候选也可能被外发 | publish 查询只筛候选 bucket,不筛 run status/audit | 仅 finished+审计通过可分发 |
| P0-08 | AIGC 按所有计划轮询,不按品类路由 | 健康、历史、时政等视频可能进入错误生产计划 | 代码明确“不区分品类”,均匀分发 | 建立可配置且可解释的分类路由 |
| P0-09 | “发布”没有真正执行发布 | 业务误以为已发布,实际只绑定了生成计划 | `publish_plan_id` 只入库,没有参与外部 API 调用 | 拆分分发/生产/发布状态并实现确认 |
| P0-10 | 外部副作用缺少端到端幂等 | 绑定失败、进程崩溃或数据库回写失败会重复创建爬取计划 | 只有 DB 回写成功后才算已处理 | 使用业务幂等键、outbox 和状态机 |
| P0-11 | AIGC 错误日志可能泄露 Token | 日志/OSS 中可能出现生产密钥 | `_post` 错误日志输出包含 `baseInfo.token` 的请求体 | 立即脱敏并轮换可能暴露的密钥 |
### 17.2 P1:高优先级正确性与稳定性问题
| ID | 问题 | 影响 | 建议 |
|---|---|---|---|
| P1-01 | 进程锁无法防止多 API worker/多副本重复调度 | 重复写库、重复模型调用、重复 AIGC 操作 | 独立 Scheduler 或分布式锁 |
| P1-02 | 需求池仅凭行数相同就跳过同步 | 行内容变化但总数不变时保留旧数据 | 比较主键+内容 hash 或直接 upsert diff |
| P1-03 | 分类树和元素只追加不更新/失效 | 分类名称、层级、父子关系长期漂移 | 引入快照版本和变更同步 |
| P1-04 | 关系边只增不删且缺少外键 | 需求池删除后留下悬空关系,数据不断膨胀 | 增加 FK/唯一键并按快照清理 |
| P1-05 | `multi_demand_pool_di` 无数据库唯一约束 | 并发执行时可产生重复业务行 | 增加 `(biz_dt,strategy,demand_id)` 唯一键 |
| P1-06 | 分类 ID 通过 `max(id)+1` 分配 | 多实例同步时可能主键冲突 | 使用自增 ID,先父后子映射或事务锁 |
| P1-07 | 一词只能挂一个分类,与 Prompt/业务目标冲突 | 多义词和跨领域需求被错误压缩 | 拆分需求词表和多对多挂靠表 |
| P1-08 | 子串匹配直接建立权威关系 | 短词误命中、语义污染、错误热度归因 | 候选召回+规则/模型确认+置信度 |
| P1-09 | 当日 `hung_word_count` 包含所有历史归类词 | 无当日数据的历史节点仍进入分级计划 | 按当日有效关系计算挂载数 |
| P1-10 | “近 7 日”查询实际覆盖 biz_dt-7 到 biz_dt,共 8 个自然日 | 后验口径偏移 | 修正为 7 个自然日并固化口径测试 |
| P1-11 | ROV/VOV 只取前 1000 行且倾向保留较高值 | 长尾缺失、后验结果乐观偏差 | 全量/分页拉取;按业务主键聚合 |
| P1-12 | 视频解析分区使用运行日昨天,不使用 biz_dt | 历史重跑读错分区,迟到数据无法补齐 | 所有分区从 biz_dt 派生 |
| P1-13 | 已存在视频详情永远跳过 | 标题、点位或解析结果更新无法自动刷新 | 保存源版本并支持变更 upsert |
| P1-14 | 无拓展候选的 S/A 需求完全不进入找片 | 高优需求因缺少拓展点而丢失 | 原需求本身作为默认搜索根 |
| P1-15 | 找片排序实现与注释不一致 | A 级高 score 可排在 S 级前 | 产品确认排序并加测试 |
| P1-16 | `running` 找片记录无租约,可能永久跳过 | 进程硬退出后任务永远不再执行 | 增加 heartbeat、超时和 attempt |
| P1-17 | 强制重跑复用旧 run_id 和旧子记录 | 两次搜索轨迹、候选和状态互相污染 | 每次重跑新 attempt,显式继承关系 |
| P1-18 | 同一 aweme_id 可在多个 run 重复分发 | AIGC 重复抓取/生产同一视频 | 建立全局视频资产和发布唯一性 |
| P1-19 | 最终文字覆盖数据库分池(已修复) | 防止未经审计的文本解析改变候选状态 | 最终报告只读数据库最终状态并由 guard 校验 |
| P1-20 | 前端“视频发现”未展示真实发现候选 | 运营无法核对主推荐、备选和发布状态 | 新增 candidate/run API 和页面 |
| P1-21 | Scheduler 默认启用且随 API 启动 | 开发、扩容或临时环境可能误触生产任务 | 生产显式开启,默认关闭 |
| P1-22 | 延迟超过 1 小时会丢失当日调度 | API 故障恢复后不会自动补跑 | 按 biz_dt 对账并自动补批次 |
| P1-23 | API 无认证并监听 `0.0.0.0` | 业务数据和 OSS 日志地址可能被未授权访问 | 接入认证、权限和网关 |
| P1-24 | 自动化测试当前无法完成收集 | 关键回归无法执行 | 修复/删除失效测试并纳入 CI |
### 17.3 P2:产品一致性与可维护性问题
| ID | 问题 | 影响 | 建议 |
|---|---|---|---|
| P2-01 | 只记录总流水线 start/end | 无法可靠查询每个子步骤的历史与重试 | 建立批次表和步骤执行表 |
| P2-02 | `scheduler_job_execution.detail` 可能写入巨大嵌套结果 | TEXT 超限后执行记录静默丢失 | 摘要字段结构化,详细结果独立存储 |
| P2-03 | Agent Prompt、模型和策略未作为业务版本落库 | 跨日等级不可重现 | 保存模型、Prompt hash、参数和工具版本 |
| P2-04 | `total_score` 将热度和维度覆盖混在一起 | 缺维节点天然低分,口径难解释 | 分离热度、覆盖率、置信度 |
| P2-05 | API 全量返回分类树和需求,缺少分页/缓存 | 数据增长后响应慢、前端内存压力大 | 子树 API、分页、ETag/缓存 |
| P2-06 | `generated_demand` 是孤立产物 | 项目存在两套“需求”概念 | 接入正式生命周期或下线 |
| P2-07 | 文档、Agent 清单、定时时间和代码不一致 | 运维和研发容易按错误流程操作 | 将本文设为主文档并持续更新 |
| P2-08 | `create_find_agent` 的 model 参数未实际生效 | 调试和灰度模型切换失效 | 尊重传入 model,并记录版本 |
| P2-09 | 依赖清单未显式声明 `requests` | 依赖传递变化时 AIGC 客户端可能无法启动 | 在项目依赖中直接声明 |
---
## 18. 已复现的测试问题
执行:
```bash
.venv/bin/python -m pytest -q
```
当前在测试收集阶段失败,未进入完整测试执行:
1. `tests/agents/demand_grade_orchestrator_agent/test_orchestrator.py` 引用已不存在的
`common.plan_builder`;
2. `tests/supply_infra/scheduler/test_grade_demand_pool.py` 引用已从当前实现移除的
`_execute_plan_tasks_with_retries` 等函数。
这说明当前“增加全流程定时任务测试”的提交与实际代码没有保持同步,不能将现有测试目录视为有效回归保障。
---
## 19. 目标流程
```mermaid
flowchart TB
BATCH["创建每日批次
冻结 biz_dt/代码/策略/模型版本"]
LOCK{"取得分布式锁?"}
SYNC["同步并校验树、需求池、后验、视频分区"]
QUALITY{"数据质量门禁通过?"}
GRADE["生成分级计划并执行"]
COVER{"计划覆盖率=100%
分级覆盖率=100%?"}
EXPAND["S/A 拓展
每条都有明确终态"]
DISCOVER["找片 attempt
租约/恢复/完成审计"]
AUDIT{"运行 finished
候选审计通过?"}
ROUTE["按需求分类路由 AIGC"]
OUTBOX["幂等 outbox
创建→绑定→生产→发布"]
VERIFY["查询外部状态并确认"]
DONE["批次完成并发布日报"]
STOP["停止下游
告警+补偿队列"]
BATCH --> LOCK
LOCK -->|否| STOP
LOCK -->|是| SYNC
SYNC --> QUALITY
QUALITY -->|否| STOP
QUALITY -->|是| GRADE
GRADE --> COVER
COVER -->|否| STOP
COVER -->|是| EXPAND
EXPAND --> DISCOVER
DISCOVER --> AUDIT
AUDIT -->|否| STOP
AUDIT -->|是| ROUTE
ROUTE --> OUTBOX
OUTBOX --> VERIFY
VERIFY -->|失败| STOP
VERIFY -->|成功| DONE
```
---
## 20. 迭代优先级
### 20.1 第一阶段:阻止错误外发
- 为总流水线增加依赖门禁;
- 修复 failed 分级组被视为完成的问题;
- 增加计划覆盖率、分级覆盖率和逐项落库校验;
- 保持找片完成守卫、数据库审计和最终状态顺序的负向回归;
- AIGC 只读取 finished 且审计通过的运行;
- 暂停无分类路由的自动分发,先切换为 dry-run 或人工确认;
- 对 AIGC 请求日志脱敏;
- 明确“分发、生产、发布”三种状态。
### 20.2 第二阶段:实现可靠重跑
- 引入每日批次、步骤状态、attempt 和分布式锁;
- 增加业务唯一键和外键;
- 修复需求池 count-skip、关系清理和树变更同步;
- 所有分区从 `biz_dt` 派生;
- 找片运行增加租约、超时和恢复;
- AIGC 外部调用改为 outbox + 幂等状态机;
- 全局 aweme 去重。
### 20.3 第三阶段:完善产品闭环
- 重构原始需求、标准词、平台需求和多挂靠关系;
- 将找片结果和 AIGC 状态接入 API/前端;
- 接入真实生产/发布结果和 ROV/VOV 回流;
- 统一或下线 `generated_demand` 旁路;
- 增加策略中心、版本比较和人工反馈。
---
## 21. 验收标准
### 21.1 每日运行
- 同一环境、同一 `biz_dt` 只产生一个有效主批次;
- 任何上游失败都不会触发错误下游外发;
- 每个步骤都有独立状态、耗时、输入量、输出量和错误摘要;
- 延迟或停机恢复后能自动识别并补跑缺失业务日。
### 21.2 数据正确性
- 需求池源数据主键和内容与 ODPS 对账一致;
- 分类树新增、更新、迁移、删除都有明确处理;
- 时间窗口自动化测试证明“7 日”恰好包含 7 个业务日;
- 当日挂载数只统计当日有效需求;
- 历史重跑读取对应 `biz_dt` 的全部分区。
### 21.3 分级
- 计划节点覆盖率 100%;
- 当日有效需求结果覆盖率 100%;
- failed 组数大于 0 时总步骤必为失败;
- 每个 Agent 批次输入都能在数据库逐条找到结果或明确失败原因;
- 相同数据、模型和策略版本重跑结果可解释、可比较。
### 21.4 找片
- 每条入选 S/A 需求都有“已完成、无候选、无数据、失败”之一;
- 无拓展点的 S/A 需求仍能以原需求搜索;
- finished 运行不存在 `pending_evaluation` 候选;
- 搜索页、证据、候选分池、审计和最终报告一致;
- 崩溃遗留 running 能在超时后自动恢复。
### 21.5 AIGC
- 100% 候选来自 finished 且审计通过的运行;
- 每个候选进入与需求分类匹配的计划;
- 同一 aweme 在同一发布策略下最多外发一次;
- 创建、绑定、生产、发布状态可分别查询;
- 外部接口和数据库任一侧重试不会产生重复计划;
- 日志中不存在明文 Token。
### 21.6 测试和发布
- `pytest` 可完整收集并通过;
- P0 路径具备单元测试和集成测试;
- CI 覆盖同步差异、失败门禁、断点重跑、并发锁、找片审计和 AIGC 幂等;
- AIGC 默认 dry-run,通过灰度和人工确认后才开启真实外发。
---
## 22. 建议运营指标
| 指标 | 建议目标 |
|---|---|
| 每日批次成功率 | ≥ 99% |
| 数据同步对账差异 | 0 |
| 分级计划覆盖率 | 100% |
| 分级结果覆盖率 | 100% |
| S/A 明确终态覆盖率 | 100% |
| 找片审计通过率 | 可按失败原因分层,不允许绕过 |
| 重复 AIGC 外发率 | 0 |
| 错误品类路由率 | 0 |
| 密钥明文日志事件 | 0 |
| P0 失败发现时延 | ≤ 5 分钟 |
---
## 23. 待业务确认
1. 每日 `biz_dt` 应使用当天还是 T-1,15:00 时上游当天分区是否已稳定;
2. S 与 A 的资源排序是否必须严格 S 优先;
3. 每天 Top 200 是固定预算,还是应按分类、等级、探索比例动态分配;
4. 没有源视频或点位的 S/A 需求应如何搜索;
5. AIGC 计划应按哪个分类层级路由,一条需求多挂靠时如何选主路由;
6. “发布”最终是完成生产计划绑定,还是必须真正上线到目标账号;
7. ROV/VOV 的有效样本门槛、归因方式和跨日衰减规则;
8. `generated_demand` 是否要成为正式平台需求,还是退出主产品;
9. API 和 OSS 日志的访问权限、数据保留周期及脱敏要求。
---
## 24. 代码导航
| 流程 | 关键代码 |
|---|---|
| Scheduler 注册 | `supply_infra/scheduler/app.py` |
| 总流水线 | `supply_infra/scheduler/jobs/run_supply_pipeline.py` |
| 全局树同步 | `supply_infra/scheduler/jobs/sync_global_tree_odps_to_mysql.py` |
| 需求池同步 | `supply_infra/scheduler/jobs/demand_pool/` |
| 树热度 | `supply_infra/scheduler/jobs/demand_pool/tree_weight.py` |
| 分级 | `supply_infra/scheduler/jobs/grade_demand_pool.py`、`agents/demand_grade_agent/` |
| 点位拓展 | `supply_infra/scheduler/jobs/expand_demand_from_video_points.py`、`agents/demand_video_expand_agent/` |
| 找片 | `supply_infra/scheduler/jobs/discover_videos_from_demands.py`、`agents/find_agent/` |
| AIGC 分发 | `supply_infra/scheduler/jobs/publish_videos_from_discovery.py`、`supply_infra/aigc/` |
| 数据模型 | `supply_infra/db/models/` |
| API | `api/app.py`、`api/services/` |
| 前端 | `web/src/views/`、`web/src/components/` |