文档性质:现状还原 + 目标产品要求 + 风险整改清单 版本:v1.0 更新日期:2026-07-24 代码基线:
dev-find/5ea5e65核心入口:supply_infra/scheduler/jobs/run_supply_pipeline.py
本文用于统一说明 SupplyAgent 当前项目的完整业务流程,重点回答:
本文以当前代码实际行为为准。仓库中的 README.md、ARCHITECTURE.md、prd/ 和
zhangbo.md 部分内容描述的是旧流程或目标蓝图,不能代替本文对现状的说明。
本次分析未触发 ODPS、外部搜索、AIGC 发布等线上副作用;运行态数据量、接口成功率和真实耗时仍需结合生产监控补充。
SupplyAgent 是一套面向内容供给的每日需求处理系统。它将多来源需求信号组织到全局分类树中, 结合先验热度和真实效果完成需求分级,再从已有视频点位中扩展搜索意图,自动寻找适合目标受众的 短视频,并把保留候选分发到 AIGC 生产计划。
当前产品形态由五部分组成:
当前代码尚未完整实现以下目标:
find_agent 最终候选及完整审计过程;| 角色/系统 | 主要诉求 | 当前交互 |
|---|---|---|
| 内容策略/运营 | 看懂需求热度、等级、原因和内容证据 | 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 平台 | 创建视频爬取计划并绑定生产计划 | 流水线末端调用 |
| 对象 | 说明 | 核心数据表 |
|---|---|---|
| 全局分类树 | 稳定的内容分类骨架 | 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 |
flowchart LR
START["API/CLI 启动 Scheduler"] --> CRON["每日 15:00<br/>Asia/Shanghai"]
CRON --> T1["① 同步 T-1 全局分类树"]
T1 --> POOL["② 同步当日策略需求池"]
POOL --> CLASSIFY["归类新词并建立匹配边"]
CLASSIFY --> METRIC["回填 ROV/VOV<br/>计算词级和树级热度"]
METRIC --> SOURCE_VIDEO["同步源视频标题与三类点位"]
SOURCE_VIDEO --> GRADE["③ 需求分级<br/>S/A/B/C/D"]
GRADE --> EXPAND["④ S/A 视频点位拓展"]
EXPAND --> FIND["⑤ Top 200 需求找片"]
FIND --> AIGC["⑥ 创建 AIGC 爬取计划<br/>绑定生产计划"]
AIGC --> API["FastAPI 查询"]
API --> WEB["Vue 需求地图/需求/视频证据"]
biz_dt 时,按 SCHEDULER_TIMEZONE 的当天生成 YYYYMMDD;biz_dt - 1 天 的 ODPS 分区;biz_dt;biz_dt,而是固定读取任务实际执行日的昨天。init_db();SCHEDULER_ENABLED=true,在 API lifespan 中启动后台 Scheduler(唯一启动入口);python -m supply_infra.scheduler.jobs.run_supply_pipeline [YYYYMMDD] 手动执行全链路。| 参数 | 当前值 |
|---|---|
| 触发时间 | 每日 15:00 |
| 时区 | 默认 Asia/Shanghai |
| Job ID | run_supply_pipeline |
| 同一 Scheduler 最大实例 | 1 |
| 合并错过的执行 | coalesce=True |
| 允许延迟 | 3600 秒 |
| 找片需求上限 | 200 |
| 找片并发 | 2 |
| 分级并发 | 5 |
| 点位拓展并发 | 5 |
threading.Lock 防止同进程重入;scheduler_job_execution;pattern_mining_element 的 T-1 分区;public_pattern_mining_category 的 T-1 分区;execution_id=401、source_type/element_type=实质 等条件。category_id;source_id 去重;INSERT IGNORE;global_tree_category;global_tree_element;本阶段实际上是一个包含七个子阶段的内部流水线。任一子阶段失败时,其余子阶段仍会继续。
输入为 ODPS dwd_multi_demand_pool_di 的 biz_dt 分区。
处理规则:
(strategy, demand_id) 去重行数;(strategy, demand_id) 做插入、删除和更新;video_list 最多保留 10 个视频;weight=0 转为 NULL;video_list、video_count 和 weight。输出为 multi_demand_pool_di 当日数据。
demand_name;demand_belong_category 已存在的名称;demand_belong_category_agent;当前表以 name 唯一,所以一个词实际上只能保存一个分类节点。
归类词.name in 需求池.demand_name,则建立关系边;输出:
demand_belong_pool_rel;demand_belong_category.video_list。biz_dt-7天 至 biz_dt 的生产效果数据;rov_diff、vov_diff;rov_diff 优先、再比较 vov_diff 的较高记录;特征值 == demand_name 精确匹配回填当日需求池。输出:
multi_demand_pool_di.real_rov_7d;multi_demand_pool_di.real_vov_7d。对 demand_belong_category 的全部有效词逐一处理:
demand_name LIKE %词% 查找匹配行;(demand_category_id, biz_dt) upsert。策略映射:
| 来源策略 | 指标 |
|---|---|
| 新热事件 | ext_pop |
| 逐月 | plat_sust_pop |
| 去年同期阳历、去年同期阴历 | plat_ly_pop |
| 当下供需gap | recent_pop |
输出为 demand_popularity_stats。
total_score;total_score。输出为 category_tree_weight。
video_list 收集视频 ID;输出:
multi_demand_video_detail;multi_demand_video_point。将当日需求分为 S/A/B/C/D,供后续点位拓展、找片和资源分配使用。
category_tree_weight.hung_word_count > 0 的节点中找待分配节点;demand_grade_plan、demand_grade_plan_group;当前每日定时任务默认使用代码自动分组,不使用 demand_grade_orchestrator_agent。
demand_grade_agent 查询:
total_score;score:
demand_grade;demand_grade_category_rel;一条需求必须同时满足:
demand_grade.video_list 非空;demand_video_expand_agent;demand_video_expansion;demand_video_expansion_run。demand_video_expansion 的需求;score 降序,再用 S/A 作为次级排序;running 或 finished 找片记录时默认跳过。video_discovery_run 并生成 run_id;video_discovery_search;video_discovery_candidate,初始为 pending_evaluation;primary:主推荐;backup:需求相关性不足但受众和分享价值较强的补充推荐;rejected:淘汰;pending_evaluation:尚未完成;finished 并输出主推荐、补充推荐、淘汰原因、搜索树和缺失证据;video_discovery_run;video_discovery_search;video_discovery_candidate;decision_bucket in (primary, backup) 的候选;aweme_id 必须非空;aigc_crawler_plan_id 的候选;video_discovery_run.status=finished。(生成ID, 发布ID) 去重;当前代码只“创建爬取计划并绑定生成计划”,没有使用 publish_plan_id 调用发布接口,
也没有验证内容生产或发布完成。因此现阶段准确名称应为“AIGC 生产计划分发”,不能视为真正发布完成。
| 接口 | 当前用途 |
|---|---|
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 日志 |
| 页面 | 当前展示 |
|---|---|
| 平台全局需求地图 | 分类树、热度、分级需求及下钻 |
| 全局分类树 | 可展开分类树和多维热度 |
| 需求归类过程 | 归类 Agent OSS 日志 |
| 需求汇总/视频发现 | 分级需求、源视频、拓展点位 |
当前前端“视频发现”页面没有读取 video_discovery_candidate,因此展示的是找片输入证据,
不是 find_agent 最终找到的主推荐、补充推荐和淘汰候选。
以下能力存在于项目中,但不属于每日总流水线:
generate_demand_agent:按单一热度维度写入 generated_demand;demand_grade_orchestrator_agent:模型统筹分组,当前定时路径已改为代码自动分组;python -m supply_infra.scheduler.jobs.grade_demand_pool --retry-failed:手动重试失败分级明细;python -m supply_infra.scheduler.jobs.* 单独跑各流水线步骤;generated_demand 当前没有接入总流水线、API 或前端,属于孤立产物。
biz_dt 建立唯一每日批次;biz_dt 派生,支持历史重跑;running 可自动恢复;primary 和 backup 是否都允许自动分发必须由产品策略明确配置;find_agent 的主推荐、补充推荐、淘汰候选和证据;| 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 | 找片完成守卫被关闭 | Agent 可在搜索、证据、评估或审计未完成时提前结束 | create_find_agent(... completion_guard=None) |
恢复并测试确定性完成守卫 |
| 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 的请求体 |
立即脱敏并轮换可能暴露的密钥 |
| 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 | 最终文字可再次覆盖数据库分池 | 未经审计的文本解析可能改变候选状态 | 最终报告从数据库生成,禁止反向覆盖 |
| 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 |
| 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 客户端可能无法启动 | 在项目依赖中直接声明 |
执行:
.venv/bin/python -m pytest -q
当前在测试收集阶段失败,未进入完整测试执行:
tests/agents/demand_grade_orchestrator_agent/test_orchestrator.py 引用已不存在的
common.plan_builder;tests/supply_infra/scheduler/test_grade_demand_pool.py 引用已从当前实现移除的
_execute_plan_tasks_with_retries 等函数。这说明当前“增加全流程定时任务测试”的提交与实际代码没有保持同步,不能将现有测试目录视为有效回归保障。
flowchart TB
BATCH["创建每日批次<br/>冻结 biz_dt/代码/策略/模型版本"]
LOCK{"取得分布式锁?"}
SYNC["同步并校验树、需求池、后验、视频分区"]
QUALITY{"数据质量门禁通过?"}
GRADE["生成分级计划并执行"]
COVER{"计划覆盖率=100%<br/>分级覆盖率=100%?"}
EXPAND["S/A 拓展<br/>每条都有明确终态"]
DISCOVER["找片 attempt<br/>租约/恢复/完成审计"]
AUDIT{"运行 finished<br/>候选审计通过?"}
ROUTE["按需求分类路由 AIGC"]
OUTBOX["幂等 outbox<br/>创建→绑定→生产→发布"]
VERIFY["查询外部状态并确认"]
DONE["批次完成并发布日报"]
STOP["停止下游<br/>告警+补偿队列"]
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
biz_dt 派生;generated_demand 旁路;biz_dt 只产生一个有效主批次;biz_dt 的全部分区。pending_evaluation 候选;pytest 可完整收集并通过;| 指标 | 建议目标 |
|---|---|
| 每日批次成功率 | ≥ 99% |
| 数据同步对账差异 | 0 |
| 分级计划覆盖率 | 100% |
| 分级结果覆盖率 | 100% |
| S/A 明确终态覆盖率 | 100% |
| 找片审计通过率 | 可按失败原因分层,不允许绕过 |
| 重复 AIGC 外发率 | 0 |
| 错误品类路由率 | 0 |
| 密钥明文日志事件 | 0 |
| P0 失败发现时延 | ≤ 5 分钟 |
biz_dt 应使用当天还是 T-1,15:00 时上游当天分区是否已稳定;backup 是否允许无人工确认直接进入 AIGC;generated_demand 是否要成为正式平台需求,还是退出主产品;| 流程 | 关键代码 |
|---|---|
| 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/ |