Procházet zdrojové kódy

增加手动触发定时任务

xueyiming před 1 týdnem
rodič
revize
5240574d49
4 změnil soubory, kde provedl 1342 přidání a 0 odebrání
  1. 844 0
      PRD.md
  2. 102 0
      api/app.py
  3. 46 0
      api/services/scheduler.py
  4. 350 0
      supply_infra/scheduler/manual_jobs.py

+ 844 - 0
PRD.md

@@ -0,0 +1,844 @@
+# 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 状态接口、日志、`jobs/` CLI |
+| 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["每日 14:30<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 需求地图/需求/视频证据"]
+```
+
+### 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 jobs/run_scheduler.py` 或 `supply-scheduler` 单独启动;
+4. 可通过 `python jobs/run_supply_pipeline.py [YYYYMMDD]` 手动执行全链路。
+
+### 6.2 调度参数
+
+| 参数 | 当前值 |
+|---|---|
+| 触发时间 | 每日 14:30 |
+| 时区 | 默认 `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`:主推荐;
+   - `backup`:需求相关性不足但受众和分享价值较强的补充推荐;
+   - `rejected`:淘汰;
+   - `pending_evaluation`:尚未完成;
+9. Agent 应将运行置为 `finished` 并输出主推荐、补充推荐、淘汰原因、搜索树和缺失证据;
+10. Agent 返回后,程序从最终文字中解析主推荐/补充推荐视频 ID,再次同步数据库分池。
+
+### 11.3 输出
+
+- `video_discovery_run`;
+- `video_discovery_search`;
+- `video_discovery_candidate`;
+- Agent 审计日志和 OSS HTML。
+
+---
+
+## 12. 阶段六:AIGC 分发
+
+### 12.1 当前入选规则
+
+- 选择当日 `decision_bucket in (primary, backup)` 的候选;
+- `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`:模型统筹分组,当前定时路径已改为代码自动分组;
+- 多个 `jobs/backfill_*`:补标题、视频列表、点位和描述;
+- `retry_failed_grade_plan_items.py`:手动重试失败分级明细;
+- 单独运行分类树权重、排名、需求池关系、视频同步、找片和 AIGC 分发;
+- 日志可视化和手动 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` 和 `backup` 是否都允许自动分发必须由产品策略明确配置;
+- 候选必须按需求分类或明确的路由规则进入正确计划;
+- 创建、绑定、生产、发布每个外部动作必须有幂等键和独立状态;
+- 必须区分“已分发、已生产、已发布、发布失败”;
+- 只有真实发布接口成功并经查询确认后,才能标记“已发布”。
+
+### 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 | 找片完成守卫被关闭 | 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` 的请求体 | 立即脱敏并轮换可能暴露的密钥 |
+
+### 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 | 最终文字可再次覆盖数据库分池 | 未经审计的文本解析可能改变候选状态 | 最终报告从数据库生成,禁止反向覆盖 |
+| 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["创建每日批次<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
+```
+
+---
+
+## 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,14:30 时上游当天分区是否已稳定;
+2. S 与 A 的资源排序是否必须严格 S 优先;
+3. `backup` 是否允许无人工确认直接进入 AIGC;
+4. 每天 Top 200 是固定预算,还是应按分类、等级、探索比例动态分配;
+5. 没有源视频或点位的 S/A 需求应如何搜索;
+6. AIGC 计划应按哪个分类层级路由,一条需求多挂靠时如何选主路由;
+7. “发布”最终是完成生产计划绑定,还是必须真正上线到目标账号;
+8. ROV/VOV 的有效样本门槛、归因方式和跨日衰减规则;
+9. `generated_demand` 是否要成为正式平台需求,还是退出主产品;
+10. API 和 OSS 日志的访问权限、数据保留周期及脱敏要求。
+
+---
+
+## 24. 代码导航
+
+| 流程 | 关键代码 |
+|---|---|
+| Scheduler 注册 | `supply_infra/scheduler/app.py` |
+| 总流水线 | `supply_infra/scheduler/jobs/run_supply_pipeline.py` |
+| 全局树同步 | `sync_global_tree_odps_to_mysql.py` |
+| 需求池内部流水线 | `sync_multi_demand_pool_odps_to_mysql.py` |
+| 树热度 | `compute_category_tree_weight.py` |
+| 分级 | `grade_demand_pool.py`、`agents/demand_grade_agent/` |
+| 点位拓展 | `expand_demand_from_video_points.py`、`agents/demand_video_expand_agent/` |
+| 找片 | `discover_videos_from_demands.py`、`agents/find_agent/` |
+| AIGC 分发 | `publish_videos_from_discovery.py`、`supply_infra/aigc/` |
+| 数据模型 | `supply_infra/db/models/` |
+| API | `api/app.py`、`api/services/` |
+| 前端 | `web/src/views/`、`web/src/components/` |

+ 102 - 0
api/app.py

@@ -7,6 +7,7 @@ from pathlib import Path
 from fastapi import FastAPI, HTTPException, Query
 from fastapi.middleware.cors import CORSMiddleware
 from fastapi.staticfiles import StaticFiles
+from pydantic import BaseModel, Field
 
 from api.services.category_tree import build_category_tree
 from api.services.demand_belong_category import list_demand_belong_categories
@@ -14,6 +15,12 @@ from api.services.demand_grade import list_demand_grades
 from api.services.demand_grade_videos import list_videos_for_demand_grade
 from api.services.demand_videos import list_videos_for_demand_belong
 from api.services.oss_logs import list_demand_belong_oss_logs
+from api.services.scheduler import (
+    get_scheduler_job_run,
+    list_triggerable_jobs,
+    run_scheduler_job,
+    run_supply_pipeline,
+)
 from api.services.video_discovery import (
     get_video_discovery_demand,
     list_video_discovery_demands,
@@ -57,6 +64,101 @@ def scheduler_status() -> dict:
     return get_scheduler_status()
 
 
+class RunPipelineBody(BaseModel):
+    biz_dt: str | None = Field(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="业务日 YYYYMMDD;省略则取当天",
+    )
+
+
+class TriggerSchedulerJobBody(BaseModel):
+    biz_dt: str | None = Field(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="业务日 YYYYMMDD;省略则取当天",
+    )
+    partition_date: str | None = Field(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="ODPS 分区日 YYYYMMDD;全局树默认取 biz_dt 前一日",
+    )
+    wait: bool = Field(
+        default=False,
+        description="是否同步等待执行完成;默认 false 表示后台执行",
+    )
+
+
+@app.post("/api/pipeline/run", status_code=202)
+def run_pipeline(
+    body: RunPipelineBody | None = None,
+    biz_dt: str | None = Query(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="业务日 YYYYMMDD(与 body 二选一,body 优先)",
+    ),
+) -> dict:
+    """
+    异步一键执行供给数据全流程,立即返回 run_id,后台串行执行:
+
+    全局树同步 → 需求池同步 → 需求分级 → 视频点位拓展 → find_agent 找片 → AIGC 发布。
+    """
+    resolved_biz_dt = body.biz_dt if body and body.biz_dt is not None else biz_dt
+    return run_supply_pipeline(biz_dt=resolved_biz_dt)
+
+
+@app.get("/api/scheduler/jobs")
+def scheduler_jobs() -> dict:
+    """Return jobs that can be triggered manually."""
+    return {"items": list_triggerable_jobs()}
+
+
+@app.post("/api/scheduler/jobs/{job_id}/run")
+def trigger_scheduler_job(
+    job_id: str,
+    body: TriggerSchedulerJobBody | None = None,
+    biz_dt: str | None = Query(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="业务日 YYYYMMDD(与 body 二选一,body 优先)",
+    ),
+    partition_date: str | None = Query(
+        default=None,
+        pattern=r"^\d{8}$",
+        description="ODPS 分区日 YYYYMMDD(与 body 二选一,body 优先)",
+    ),
+    wait: bool = Query(
+        default=False,
+        description="是否同步等待执行完成",
+    ),
+) -> dict:
+    """Manually trigger a scheduler job."""
+    resolved_biz_dt = body.biz_dt if body and body.biz_dt is not None else biz_dt
+    resolved_partition_date = (
+        body.partition_date if body and body.partition_date is not None else partition_date
+    )
+    resolved_wait = body.wait if body is not None else wait
+
+    try:
+        return run_scheduler_job(
+            job_id,
+            biz_dt=resolved_biz_dt,
+            partition_date=resolved_partition_date,
+            wait=resolved_wait,
+        )
+    except KeyError:
+        raise HTTPException(status_code=404, detail=f"scheduler job not found: {job_id}") from None
+
+
+@app.get("/api/scheduler/runs/{run_id}")
+def scheduler_job_run(run_id: str) -> dict:
+    """Return manual job run status (in-memory; cleared on process restart)."""
+    result = get_scheduler_job_run(run_id)
+    if result is None:
+        raise HTTPException(status_code=404, detail="scheduler run not found")
+    return result
+
+
 @app.get("/api/category-tree")
 def category_tree(
     biz_dt: str | None = Query(

+ 46 - 0
api/services/scheduler.py

@@ -0,0 +1,46 @@
+"""Scheduler API helpers."""
+from __future__ import annotations
+
+from typing import Any
+
+from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID
+from supply_infra.scheduler.manual_jobs import (
+    get_manual_run,
+    list_manual_jobs,
+    trigger_manual_job,
+)
+
+
+def list_triggerable_jobs() -> list[dict[str, Any]]:
+    return list_manual_jobs()
+
+
+def run_scheduler_job(
+    job_id: str,
+    *,
+    biz_dt: str | None = None,
+    partition_date: str | None = None,
+    wait: bool = False,
+) -> dict[str, Any]:
+    return trigger_manual_job(
+        job_id,
+        biz_dt=biz_dt,
+        partition_date=partition_date,
+        wait=wait,
+    )
+
+
+def get_scheduler_job_run(run_id: str) -> dict[str, Any] | None:
+    return get_manual_run(run_id)
+
+
+def run_supply_pipeline(
+    *,
+    biz_dt: str | None = None,
+) -> dict[str, Any]:
+    """异步执行供给数据全流程,立即返回 run_id。"""
+    return run_scheduler_job(
+        SUPPLY_PIPELINE_JOB_ID,
+        biz_dt=biz_dt,
+        wait=False,
+    )

+ 350 - 0
supply_infra/scheduler/manual_jobs.py

@@ -0,0 +1,350 @@
+"""手动触发定时任务:任务注册表、后台执行与运行状态查询。"""
+from __future__ import annotations
+
+import logging
+import threading
+import uuid
+from collections.abc import Callable
+from dataclasses import dataclass
+from datetime import datetime, timedelta
+from typing import Any
+from zoneinfo import ZoneInfo
+
+from supply_infra.config import get_infra_settings
+from supply_infra.scheduler.constants import (
+    SUPPLY_PIPELINE_JOB_ID,
+    SUPPLY_PIPELINE_JOB_NAME,
+)
+
+logger = logging.getLogger(__name__)
+
+_RUNS: dict[str, dict[str, Any]] = {}
+_RUNS_LOCK = threading.Lock()
+
+
+@dataclass(frozen=True)
+class ManualJobSpec:
+    job_id: str
+    name: str
+    description: str
+    runner: Callable[..., dict[str, Any]]
+    accepts_biz_dt: bool = False
+    accepts_partition_date: bool = False
+
+
+def _resolve_biz_dt(biz_dt: str | None) -> str:
+    if biz_dt:
+        return biz_dt
+    timezone = ZoneInfo(get_infra_settings().scheduler_timezone)
+    return datetime.now(timezone).strftime("%Y%m%d")
+
+
+def _resolve_tree_partition(biz_dt: str | None, partition_date: str | None) -> str:
+    if partition_date:
+        return partition_date
+    resolved_biz_dt = _resolve_biz_dt(biz_dt)
+    return (
+        datetime.strptime(resolved_biz_dt, "%Y%m%d") - timedelta(days=1)
+    ).strftime("%Y%m%d")
+
+
+def _run_supply_pipeline_job(
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> dict[str, Any]:
+    from supply_infra.scheduler.jobs.run_supply_pipeline import run_supply_pipeline
+
+    del partition_date
+    return run_supply_pipeline(biz_dt=biz_dt)
+
+
+def _run_global_tree_job(
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> dict[str, Any]:
+    from supply_infra.scheduler.jobs.sync_global_tree_odps_to_mysql import (
+        sync_global_tree_odps_to_mysql,
+    )
+
+    return sync_global_tree_odps_to_mysql(
+        partition_date=_resolve_tree_partition(biz_dt, partition_date),
+    )
+
+
+def _run_demand_pool_job(
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> dict[str, Any]:
+    from supply_infra.scheduler.jobs.sync_multi_demand_pool_odps_to_mysql import (
+        sync_multi_demand_pool_odps_to_mysql,
+    )
+
+    return sync_multi_demand_pool_odps_to_mysql(
+        partition_date=partition_date or _resolve_biz_dt(biz_dt),
+    )
+
+
+def _run_biz_dt_job(
+    import_path: str,
+    fn_name: str,
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> dict[str, Any]:
+    del partition_date
+    module = __import__(import_path, fromlist=[fn_name])
+    fn = getattr(module, fn_name)
+    return fn(biz_dt)
+
+
+def _run_no_arg_job(
+    import_path: str,
+    fn_name: str,
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> dict[str, Any]:
+    del biz_dt, partition_date
+    module = __import__(import_path, fromlist=[fn_name])
+    fn = getattr(module, fn_name)
+    return fn()
+
+
+MANUAL_JOBS: dict[str, ManualJobSpec] = {
+    SUPPLY_PIPELINE_JOB_ID: ManualJobSpec(
+        job_id=SUPPLY_PIPELINE_JOB_ID,
+        name=SUPPLY_PIPELINE_JOB_NAME,
+        description="串行执行:全局树 → 需求池 → 分级 → 拓展 → find_agent → AIGC 发布",
+        runner=_run_supply_pipeline_job,
+        accepts_biz_dt=True,
+    ),
+    "sync_global_tree_odps_to_mysql": ManualJobSpec(
+        job_id="sync_global_tree_odps_to_mysql",
+        name="全局树 ODPS 同步",
+        description="从 ODPS 同步 global_tree_category / global_tree_element 到 MySQL",
+        runner=_run_global_tree_job,
+        accepts_biz_dt=True,
+        accepts_partition_date=True,
+    ),
+    "sync_multi_demand_pool_odps_to_mysql": ManualJobSpec(
+        job_id="sync_multi_demand_pool_odps_to_mysql",
+        name="策略需求池 ODPS 同步",
+        description="同步 multi_demand_pool_di,并执行归属分类、热度统计、树权重与视频同步",
+        runner=_run_demand_pool_job,
+        accepts_biz_dt=True,
+        accepts_partition_date=True,
+    ),
+    "sync_demand_belong_pool_rel": ManualJobSpec(
+        job_id="sync_demand_belong_pool_rel",
+        name="需求词与需求池匹配",
+        description="补充 demand_belong_pool_rel 并回填 demand_belong_category.video_list",
+        runner=lambda **kwargs: _run_no_arg_job(
+            "supply_infra.scheduler.jobs.sync_demand_belong_pool_rel",
+            "sync_demand_belong_pool_rel",
+            **kwargs,
+        ),
+    ),
+    "grade_demand_pool": ManualJobSpec(
+        job_id="grade_demand_pool",
+        name="需求池分级",
+        description="对指定业务日的需求词执行分级评估",
+        runner=lambda **kwargs: _run_biz_dt_job(
+            "supply_infra.scheduler.jobs.grade_demand_pool",
+            "grade_demand_pool",
+            **kwargs,
+        ),
+        accepts_biz_dt=True,
+    ),
+    "expand_demand_from_video_points": ManualJobSpec(
+        job_id="expand_demand_from_video_points",
+        name="视频点位拓展",
+        description="对 S/A 需求执行视频点位拓展判断",
+        runner=lambda **kwargs: _run_biz_dt_job(
+            "supply_infra.scheduler.jobs.expand_demand_from_video_points",
+            "expand_demand_from_video_points",
+            **kwargs,
+        ),
+        accepts_biz_dt=True,
+    ),
+    "discover_videos_from_demands": ManualJobSpec(
+        job_id="discover_videos_from_demands",
+        name="find_agent 视频发现",
+        description="对 S/A 需求拓展点位调用 find_agent 找片",
+        runner=lambda **kwargs: _run_biz_dt_job(
+            "supply_infra.scheduler.jobs.discover_videos_from_demands",
+            "discover_videos_from_demands",
+            **kwargs,
+        ),
+        accepts_biz_dt=True,
+    ),
+    "publish_videos_from_discovery": ManualJobSpec(
+        job_id="publish_videos_from_discovery",
+        name="AIGC 视频发布",
+        description="将 find_agent 发现结果发布到 AIGC 平台",
+        runner=lambda **kwargs: _run_biz_dt_job(
+            "supply_infra.scheduler.jobs.publish_videos_from_discovery",
+            "publish_videos_from_discovery",
+            **kwargs,
+        ),
+        accepts_biz_dt=True,
+    ),
+    "sync_multi_demand_videos": ManualJobSpec(
+        job_id="sync_multi_demand_videos",
+        name="需求池视频详情同步",
+        description="从 TikHub 拉取 multi_demand_pool_di 关联视频标题与详情",
+        runner=lambda **kwargs: _run_no_arg_job(
+            "supply_infra.scheduler.jobs.sync_multi_demand_videos",
+            "sync_multi_demand_videos",
+            **kwargs,
+        ),
+    ),
+}
+
+
+def list_manual_jobs() -> list[dict[str, Any]]:
+    """返回可手动触发的任务列表。"""
+    return [
+        {
+            "id": spec.job_id,
+            "name": spec.name,
+            "description": spec.description,
+            "accepts_biz_dt": spec.accepts_biz_dt,
+            "accepts_partition_date": spec.accepts_partition_date,
+        }
+        for spec in MANUAL_JOBS.values()
+    ]
+
+
+def get_manual_job(job_id: str) -> ManualJobSpec | None:
+    return MANUAL_JOBS.get(job_id)
+
+
+def get_manual_run(run_id: str) -> dict[str, Any] | None:
+    with _RUNS_LOCK:
+        payload = _RUNS.get(run_id)
+        return dict(payload) if payload else None
+
+
+def _set_run_state(run_id: str, payload: dict[str, Any]) -> None:
+    with _RUNS_LOCK:
+        _RUNS[run_id] = payload
+
+
+def _execute_job(
+    run_id: str,
+    spec: ManualJobSpec,
+    *,
+    biz_dt: str | None,
+    partition_date: str | None,
+) -> None:
+    started_at = datetime.now().isoformat()
+    _set_run_state(
+        run_id,
+        {
+            "run_id": run_id,
+            "job_id": spec.job_id,
+            "job_name": spec.name,
+            "status": "running",
+            "biz_dt": biz_dt,
+            "partition_date": partition_date,
+            "started_at": started_at,
+        },
+    )
+    try:
+        result = spec.runner(biz_dt=biz_dt, partition_date=partition_date)
+        success = not (isinstance(result, dict) and result.get("success") is False)
+        finished_at = datetime.now().isoformat()
+        _set_run_state(
+            run_id,
+            {
+                "run_id": run_id,
+                "job_id": spec.job_id,
+                "job_name": spec.name,
+                "status": "finished" if success else "failed",
+                "biz_dt": biz_dt,
+                "partition_date": partition_date,
+                "started_at": started_at,
+                "finished_at": finished_at,
+                "result": result,
+            },
+        )
+        logger.info(
+            "Manual job finished: job_id=%s run_id=%s status=%s",
+            spec.job_id,
+            run_id,
+            "finished" if success else "failed",
+        )
+    except Exception as exc:
+        finished_at = datetime.now().isoformat()
+        logger.exception("Manual job failed: job_id=%s run_id=%s", spec.job_id, run_id)
+        _set_run_state(
+            run_id,
+            {
+                "run_id": run_id,
+                "job_id": spec.job_id,
+                "job_name": spec.name,
+                "status": "failed",
+                "biz_dt": biz_dt,
+                "partition_date": partition_date,
+                "started_at": started_at,
+                "finished_at": finished_at,
+                "error": str(exc),
+            },
+        )
+
+
+def trigger_manual_job(
+    job_id: str,
+    *,
+    biz_dt: str | None = None,
+    partition_date: str | None = None,
+    wait: bool = False,
+) -> dict[str, Any]:
+    """
+    手动触发任务。
+
+    wait=False 时在后台线程执行并立即返回 run_id;
+    wait=True 时同步执行并返回完整结果。
+    """
+    spec = get_manual_job(job_id)
+    if spec is None:
+        raise KeyError(job_id)
+
+    if spec.accepts_biz_dt and biz_dt is None:
+        biz_dt = _resolve_biz_dt(None)
+    if spec.job_id == "sync_global_tree_odps_to_mysql" and partition_date is None:
+        partition_date = _resolve_tree_partition(biz_dt, None)
+    if spec.job_id == "sync_multi_demand_pool_odps_to_mysql" and partition_date is None:
+        partition_date = biz_dt or _resolve_biz_dt(None)
+
+    run_id = str(uuid.uuid4())
+    if wait:
+        _execute_job(run_id, spec, biz_dt=biz_dt, partition_date=partition_date)
+        payload = get_manual_run(run_id)
+        assert payload is not None
+        return payload
+
+    thread = threading.Thread(
+        target=_execute_job,
+        kwargs={
+            "run_id": run_id,
+            "spec": spec,
+            "biz_dt": biz_dt,
+            "partition_date": partition_date,
+        },
+        name=f"manual-job-{job_id}-{run_id[:8]}",
+        daemon=True,
+    )
+    thread.start()
+    return {
+        "accepted": True,
+        "run_id": run_id,
+        "job_id": spec.job_id,
+        "job_name": spec.name,
+        "status": "running",
+        "biz_dt": biz_dt,
+        "partition_date": partition_date,
+    }