xueyiming 2 주 전
부모
커밋
bc1dc015f8
37개의 변경된 파일1454개의 추가작업 그리고 12개의 파일을 삭제
  1. 22 0
      .env.example
  2. 143 0
      ARCHITECTURE.md
  3. 27 0
      agents/README.md
  4. 1 0
      agents/__init__.py
  5. 8 0
      agents/demand_belong_category_agent/__init__.py
  6. 41 0
      agents/demand_belong_category_agent/agent.py
  7. 26 0
      agents/demand_belong_category_agent/run.py
  8. 26 0
      agents/demand_belong_category_agent/tools/__init__.py
  9. 8 0
      agents/find_agent/__init__.py
  10. 41 0
      agents/find_agent/agent.py
  11. 26 0
      agents/find_agent/run.py
  12. 30 0
      agents/find_agent/tools/__init__.py
  13. 226 0
      agents/find_agent/tools/douyin_search.py
  14. 173 0
      agents/find_agent/tools/qwen_video_analyze.py
  15. 8 0
      jobs/init_db.py
  16. 14 0
      jobs/run_scheduler.py
  17. 9 1
      pyproject.toml
  18. 18 0
      requirements.txt
  19. 21 11
      supply_agent/agent/core.py
  20. 1 0
      supply_infra/__init__.py
  21. 54 0
      supply_infra/config.py
  22. 6 0
      supply_infra/db/__init__.py
  23. 26 0
      supply_infra/db/base.py
  24. 6 0
      supply_infra/db/models/__init__.py
  25. 32 0
      supply_infra/db/models/global_tree_category.py
  26. 34 0
      supply_infra/db/models/global_tree_element.py
  27. 11 0
      supply_infra/db/repositories/__init__.py
  28. 34 0
      supply_infra/db/repositories/base.py
  29. 27 0
      supply_infra/db/repositories/global_tree_category_repo.py
  30. 27 0
      supply_infra/db/repositories/global_tree_element_repo.py
  31. 52 0
      supply_infra/db/session.py
  32. 5 0
      supply_infra/odps/__init__.py
  33. 96 0
      supply_infra/odps/client.py
  34. 5 0
      supply_infra/scheduler/__init__.py
  35. 47 0
      supply_infra/scheduler/app.py
  36. 1 0
      supply_infra/scheduler/jobs/__init__.py
  37. 122 0
      supply_infra/scheduler/jobs/sync_global_tree_odps_to_mysql.py

+ 22 - 0
.env.example

@@ -15,3 +15,25 @@ SKILLS_DIR=skills
 # Logging
 # Logging
 LOG_ENABLED=true
 LOG_ENABLED=true
 LOGS_DIR=logs
 LOGS_DIR=logs
+
+# 百炼 API Key
+DASHSCOPE_API_KEY=sk-ws-...
+
+# MySQL
+MYSQL_HOST=127.0.0.1
+MYSQL_PORT=3306
+MYSQL_USER=root
+MYSQL_PASSWORD=
+MYSQL_DATABASE=supply_agent
+MYSQL_POOL_SIZE=5
+MYSQL_ECHO=false
+
+# ODPS (MaxCompute)
+ODPS_ACCESS_ID=
+ODPS_ACCESS_KEY=
+ODPS_PROJECT=
+ODPS_ENDPOINT=https://service.cn.maxcompute.aliyun.com/api
+
+# Scheduler
+SCHEDULER_ENABLED=true
+SCHEDULER_TIMEZONE=Asia/Shanghai

+ 143 - 0
ARCHITECTURE.md

@@ -0,0 +1,143 @@
+# SupplyAgent — project architecture
+
+## Directory layout
+
+```
+SupplyAgent/
+│
+├── supply_agent/                  # 核心 Agent 框架(通用,不含业务)
+│   ├── agent/                     #   Agent 主类 + ReAct 循环
+│   ├── llm/                       #   OpenRouter LLM 客户端
+│   ├── tools/                     #   @tool 装饰器 + 工具注册表
+│   ├── skills/                    #   Skill 加载器
+│   ├── logging/                   #   运行日志
+│   └── config.py                  #   Agent 配置
+│
+├── supply_infra/                  # 共享基础设施(所有 Agent / 定时任务共用)
+│   ├── config.py                  #   MySQL / ODPS / Scheduler 配置
+│   ├── db/
+│   │   ├── base.py                #   SQLAlchemy Base + TimestampMixin
+│   │   ├── session.py             #   Engine + get_session()
+│   │   ├── models/                #   ORM 实体类(一张表一个文件)
+│   │   │   └── video_content.py
+│   │   └── repositories/          #   数据访问层(CRUD / upsert)
+│   │       ├── base.py
+│   │       └── video_content_repo.py
+│   ├── odps/
+│   │   └── client.py              #   ODPS 查询封装
+│   ├── scheduler/
+│   │   ├── app.py                 #   APScheduler 调度器
+│   │   └── jobs/                  #   定时任务定义
+│   │       └── sync_odps_to_mysql.py
+│   └── tools/
+│       └── db_tools.py            #   共享 MySQL 读写工具(供 Agent 调用)
+│
+├── agents/                        # 所有业务 Agent(一个 Agent 一个子目录)
+│   ├── find_agent/                #   第一个 Agent:内容发现
+│   │   ├── agent.py               #   create_find_agent() 工厂
+│   │   ├── run.py                 #   交互式运行入口
+│   │   └── tools/                 #   Agent 专属工具
+│   │       ├── douyin_search.py
+│   │       └── qwen_video_analyze.py
+│   └── <your_next_agent>/         #   新增 Agent 按此结构复制
+│       ├── agent.py
+│       ├── run.py
+│       └── tools/
+│
+├── skills/                        # 全局共享 Skills(SKILL.md)
+├── jobs/                          # CLI 入口
+│   ├── run_scheduler.py           #   启动定时任务
+│   └── init_db.py                 #   初始化数据库表
+├── examples/                      # 框架使用示例
+├── tests/
+├── logs/                          # 运行日志(自动生成)
+├── requirements.txt
+└── pyproject.toml
+```
+
+## Layer responsibilities
+
+| 层 | 包 | 职责 | 谁用 |
+|----|-----|------|------|
+| 框架层 | `supply_agent` | LLM 调用、工具/技能机制、ReAct 循环 | 所有 Agent |
+| 基础设施层 | `supply_infra` | DB ORM、ODPS、定时任务、共享 DB 工具 | 所有 Agent + 定时任务 |
+| 业务层 | `agents/*` | 具体 Agent 逻辑、专属工具、工厂函数 | 业务开发 |
+
+## Data flow
+
+```
+                    ┌─────────────┐
+                    │   ODPS      │
+                    └──────┬──────┘
+                           │ 定时任务 (每天 02:00)
+                           ▼
+┌──────────┐    ┌─────────────────────┐    ┌──────────┐
+│  Agent   │───▶│  Repository (ORM)   │◀───│  Agent   │
+│  Tools   │    │  video_content 表    │    │  Tools   │
+└──────────┘    └─────────────────────┘    └──────────┘
+                           ▲
+                           │ upsert / query
+                    ┌──────┴──────┐
+                    │   MySQL     │
+                    └─────────────┘
+```
+
+**关键原则**:Agent Tools 和定时任务**不直接写 SQL**,统一通过 `Repository` 操作 ORM 实体。
+
+## How to add a new Agent
+
+```bash
+agents/
+└── report_agent/              # 1. 新建目录
+    ├── agent.py               # 2. create_report_agent() 工厂
+    ├── run.py                   # 3. 运行入口
+    └── tools/                   # 4. 专属工具
+        └── __init__.py
+```
+
+```python
+# agents/report_agent/agent.py
+from supply_agent import Agent
+from agents.report_agent.tools import register_all_tools
+from supply_infra.tools.db_tools import query_video_content  # 可复用共享工具
+
+def create_report_agent(**kwargs) -> Agent:
+    agent = Agent(system_prompt="你是报告生成助手...", **kwargs)
+    register_all_tools(agent.tools)
+    agent.tools.from_decorated(query_video_content)  # 复用基础设施
+    return agent
+```
+
+## How to add a new ORM entity
+
+```bash
+supply_infra/db/
+├── models/new_table.py          # 1. 定义实体类
+└── repositories/new_table_repo.py  # 2. 定义 Repository
+```
+
+然后在 `models/__init__.py` 中导出,供 `init_db()` 自动建表。
+
+## How to add a new scheduled job
+
+```bash
+supply_infra/scheduler/jobs/new_job.py   # 1. 定义任务函数
+```
+
+```python
+# supply_infra/scheduler/app.py 中注册
+scheduler.add_job(new_job, trigger=CronTrigger(hour=3), id="new_job")
+```
+
+## CLI commands
+
+```bash
+# 初始化数据库表
+python jobs/init_db.py
+
+# 启动定时任务
+python jobs/run_scheduler.py
+
+# 运行 find_agent
+python agents/find_agent/run.py
+```

+ 27 - 0
agents/README.md

@@ -0,0 +1,27 @@
+# 业务 Agent 目录
+
+每个 Agent 一个子目录,结构统一:
+
+```
+agents/<agent_name>/
+├── agent.py          # create_xxx_agent() 工厂函数(必须)
+├── run.py            # 交互式运行入口(推荐)
+├── tools/            # Agent 专属工具
+│   ├── __init__.py   # ALL_TOOLS + register_all_tools()
+│   └── my_tool.py
+└── skills/           # Agent 专属 Skills(可选,也可用全局 skills/)
+    └── my-skill/
+        └── SKILL.md
+```
+
+## 已有 Agent
+
+| Agent | 目录 | 说明 |
+|-------|------|------|
+| find_agent | `agents/find_agent/` | 抖音搜索 + 视频解析 + 内容入库 |
+
+## 新增 Agent 模板
+
+复制 `find_agent/` 目录,改名后修改 `agent.py` 中的 system_prompt 和工具注册即可。
+
+共享能力(MySQL 读写、ODPS 查询)从 `supply_infra` 引入,不要重复实现。

+ 1 - 0
agents/__init__.py

@@ -0,0 +1 @@
+"""All business agents live under this package."""

+ 8 - 0
agents/demand_belong_category_agent/__init__.py

@@ -0,0 +1,8 @@
+"""
+find_agent — 内容发现 Agent
+
+职责:抖音搜索、视频解析、内容入库。
+"""
+from agents.find_agent.agent import create_find_agent
+
+__all__ = ["create_find_agent"]

+ 41 - 0
agents/demand_belong_category_agent/agent.py

@@ -0,0 +1,41 @@
+"""
+find_agent 工厂 — 组装 Agent 实例。
+
+每个业务 Agent 都应提供 create_xxx_agent() 工厂函数,
+统一注册本 Agent 的工具 + 共享基础设施工具。
+"""
+from __future__ import annotations
+
+from supply_agent import Agent
+from supply_agent.config import Settings
+from agents.find_agent.tools import register_all_tools
+
+DEMAND_BELONG_CATEGORY_AGENT_SYSTEM_PROMPT = """\
+你是内容发现助手(find_agent),帮助用户搜索和分析短视频内容。
+
+## 工作流程
+1. 使用 douyin_search 搜索抖音视频
+2. 使用 qwen_video_analyze 解析视频内容
+3. 使用 save_video_content 将结果存入数据库
+4. 使用 query_video_content 查询历史数据
+
+请按步骤执行,给出清晰的分析报告。
+"""
+
+
+def create_find_agent(
+    settings: Settings | None = None,
+    *,
+    model: str | None = None,
+) -> Agent:
+    """创建 find_agent 实例,注册所有相关工具。"""
+    agent = Agent(
+        settings=settings,
+        model=model,
+        system_prompt=DEMAND_BELONG_CATEGORY_AGENT_SYSTEM_PROMPT,
+    )
+
+    # 本 Agent 专属工具
+    register_all_tools(agent.tools)
+
+    return agent

+ 26 - 0
agents/demand_belong_category_agent/run.py

@@ -0,0 +1,26 @@
+#!/usr/bin/env python3
+"""Run find_agent interactively."""
+
+from agents.demand_belong_category_agent import create_find_agent
+
+
+def main() -> None:
+    agent = create_find_agent()
+    print(f"demand_belong_category_agent ready | model={agent.model}")
+    print(f"tools: {agent.tools.list_tools()}")
+    print()
+
+    while True:
+        try:
+            user_input = input("You> ").strip()
+        except (EOFError, KeyboardInterrupt):
+            print("\nBye.")
+            break
+        if not user_input or user_input.lower() in ("exit", "quit", "q"):
+            break
+        result = agent.run(user_input)
+        print(f"\nAgent> {result.content}\n")
+
+
+if __name__ == "__main__":
+    main()

+ 26 - 0
agents/demand_belong_category_agent/tools/__init__.py

@@ -0,0 +1,26 @@
+"""
+find_agent 工具包
+
+在此注册本 Agent 专属的工具。
+"""
+from __future__ import annotations
+
+from collections.abc import Callable
+from typing import Any
+
+from agents.find_agent.tools.douyin_search import douyin_search
+from agents.find_agent.tools.qwen_video_analyze import qwen_video_analyze
+from supply_agent.tools.registry import ToolRegistry
+
+ALL_TOOLS: list[Callable[..., Any]] = [
+]
+
+__all__ = [
+    "ALL_TOOLS",
+    "register_all_tools",
+]
+
+
+def register_all_tools(registry: ToolRegistry) -> ToolRegistry:
+    """将 find_agent 包内的所有工具注册到 ToolRegistry。"""
+    return registry.from_decorated(*ALL_TOOLS)

+ 8 - 0
agents/find_agent/__init__.py

@@ -0,0 +1,8 @@
+"""
+find_agent — 内容发现 Agent
+
+职责:抖音搜索、视频解析、内容入库。
+"""
+from agents.find_agent.agent import create_find_agent
+
+__all__ = ["create_find_agent"]

+ 41 - 0
agents/find_agent/agent.py

@@ -0,0 +1,41 @@
+"""
+find_agent 工厂 — 组装 Agent 实例。
+
+每个业务 Agent 都应提供 create_xxx_agent() 工厂函数,
+统一注册本 Agent 的工具 + 共享基础设施工具。
+"""
+from __future__ import annotations
+
+from supply_agent import Agent
+from supply_agent.config import Settings
+from agents.find_agent.tools import register_all_tools
+
+FIND_AGENT_SYSTEM_PROMPT = """\
+你是内容发现助手(find_agent),帮助用户搜索和分析短视频内容。
+
+## 工作流程
+1. 使用 douyin_search 搜索抖音视频
+2. 使用 qwen_video_analyze 解析视频内容
+3. 使用 save_video_content 将结果存入数据库
+4. 使用 query_video_content 查询历史数据
+
+请按步骤执行,给出清晰的分析报告。
+"""
+
+
+def create_find_agent(
+    settings: Settings | None = None,
+    *,
+    model: str | None = None,
+) -> Agent:
+    """创建 find_agent 实例,注册所有相关工具。"""
+    agent = Agent(
+        settings=settings,
+        model=model,
+        system_prompt=FIND_AGENT_SYSTEM_PROMPT,
+    )
+
+    # 本 Agent 专属工具
+    register_all_tools(agent.tools)
+
+    return agent

+ 26 - 0
agents/find_agent/run.py

@@ -0,0 +1,26 @@
+#!/usr/bin/env python3
+"""Run find_agent interactively."""
+
+from agents.find_agent import create_find_agent
+
+
+def main() -> None:
+    agent = create_find_agent()
+    print(f"find_agent ready | model={agent.model}")
+    print(f"tools: {agent.tools.list_tools()}")
+    print()
+
+    while True:
+        try:
+            user_input = input("You> ").strip()
+        except (EOFError, KeyboardInterrupt):
+            print("\nBye.")
+            break
+        if not user_input or user_input.lower() in ("exit", "quit", "q"):
+            break
+        result = agent.run(user_input)
+        print(f"\nAgent> {result.content}\n")
+
+
+if __name__ == "__main__":
+    main()

+ 30 - 0
agents/find_agent/tools/__init__.py

@@ -0,0 +1,30 @@
+"""
+find_agent 工具包
+
+在此注册本 Agent 专属的工具。
+"""
+from __future__ import annotations
+
+from collections.abc import Callable
+from typing import Any
+
+from agents.find_agent.tools.douyin_search import douyin_search
+from agents.find_agent.tools.qwen_video_analyze import qwen_video_analyze
+from supply_agent.tools.registry import ToolRegistry
+
+ALL_TOOLS: list[Callable[..., Any]] = [
+    douyin_search,
+    qwen_video_analyze,
+]
+
+__all__ = [
+    "ALL_TOOLS",
+    "douyin_search",
+    "qwen_video_analyze",
+    "register_all_tools",
+]
+
+
+def register_all_tools(registry: ToolRegistry) -> ToolRegistry:
+    """将 find_agent 包内的所有工具注册到 ToolRegistry。"""
+    return registry.from_decorated(*ALL_TOOLS)

+ 226 - 0
agents/find_agent/tools/douyin_search.py

@@ -0,0 +1,226 @@
+"""
+抖音关键词搜索工具
+
+调用内部爬虫服务进行抖音关键词搜索。
+"""
+from __future__ import annotations
+
+import asyncio
+import json
+import logging
+import time
+from typing import Any, Optional
+
+import httpx
+
+from supply_agent.tools import tool
+
+logger = logging.getLogger(__name__)
+
+_MIN_REQUEST_INTERVAL_SECONDS = 10.1
+_rate_limit_lock = asyncio.Lock()
+_last_request_monotonic: float = 0.0
+
+# API 基础配置
+DOUYIN_SEARCH_API = "http://crawapi.piaoquantv.com/crawler/dou_yin/keyword"
+DEFAULT_TIMEOUT = 60.0
+DOUYIN_ACCOUNT_ID = "771431222"
+
+
+def _build_search_results(items: list[dict[str, Any]]) -> list[dict[str, Any]]:
+    """将 API 原始条目转换为结构化搜索结果。"""
+    results = []
+    for item in items:
+        author = item.get("author", {}) if isinstance(item.get("author"), dict) else {}
+        stats = item.get("statistics", {}) if isinstance(item.get("statistics"), dict) else {}
+        aweme_id = item.get("aweme_id", "")
+        results.append(
+            {
+                "aweme_id": aweme_id,
+                "desc": (item.get("desc") or item.get("item_title") or "无标题")[:100],
+                "url": f"https://www.douyin.com/video/{aweme_id}" if aweme_id else "",
+                "author": {
+                    "nickname": author.get("nickname", "未知作者"),
+                    "sec_uid": author.get("sec_uid", ""),
+                },
+                "statistics": {
+                    "digg_count": stats.get("digg_count", 0),
+                    "comment_count": stats.get("comment_count", 0),
+                    "share_count": stats.get("share_count", 0),
+                },
+            }
+        )
+    return results
+
+
+def _build_output_summary(
+    keyword: str,
+    items: list[dict[str, Any]],
+    has_more: bool,
+    cursor_value: str,
+) -> str:
+    """生成给 LLM 阅读的文本摘要。"""
+    lines = [f"搜索关键词「{keyword}」"]
+    lines.append(
+        f"找到 {len(items)} 条结果"
+        + (f",还有更多(cursor={cursor_value})" if has_more else "")
+    )
+    lines.append("")
+
+    for i, item in enumerate(items, 1):
+        aweme_id = item.get("aweme_id", "unknown")
+        desc = (item.get("desc") or item.get("item_title") or "无标题")[:50]
+
+        author = item.get("author", {}) if isinstance(item.get("author"), dict) else {}
+        author_name = author.get("nickname", "未知作者")
+        author_id = author.get("sec_uid", "")
+
+        stats = item.get("statistics", {}) if isinstance(item.get("statistics"), dict) else {}
+        digg_count = stats.get("digg_count", 0)
+        comment_count = stats.get("comment_count", 0)
+        share_count = stats.get("share_count", 0)
+
+        lines.append(f"{i}. {desc}")
+        lines.append(f"   ID: {aweme_id}")
+        lines.append(f"   链接: https://www.douyin.com/video/{aweme_id}")
+        lines.append(f"   作者: {author_name}")
+        lines.append(f"   sec_uid: {author_id}")
+        lines.append(f"   数据: 点赞 {digg_count:,} | 评论 {comment_count:,} | 分享 {share_count:,}")
+        lines.append("")
+
+    return "\n".join(lines)
+
+
+def _success_result(
+    keyword: str,
+    data: dict[str, Any],
+    items: list[dict[str, Any]],
+    has_more: bool,
+    cursor_value: str,
+    duration_ms: int,
+) -> str:
+    """构建成功时的 JSON 字符串返回值。"""
+    search_results = _build_search_results(items)
+    payload = {
+        "title": f"抖音搜索: {keyword}",
+        "output": _build_output_summary(keyword, items, has_more, cursor_value),
+        "keyword": keyword,
+        "results_count": len(items),
+        "has_more": has_more,
+        "next_cursor": cursor_value,
+        "search_results": search_results,
+        "duration_ms": duration_ms,
+    }
+    return json.dumps(payload, ensure_ascii=False)
+
+
+def _error_result(error: str, *, title: str = "抖音搜索失败") -> str:
+    """构建失败时的 JSON 字符串返回值。"""
+    return json.dumps({"error": error, "title": title}, ensure_ascii=False)
+
+
+@tool
+async def douyin_search(
+    keyword: str,
+    content_type: str = "视频",
+    sort_type: str = "综合排序",
+    publish_time: str = "不限",
+    cursor: str = "0",
+    account_id: str = DOUYIN_ACCOUNT_ID,
+    timeout: Optional[float] = None,
+) -> str:
+    """
+    抖音关键词搜索
+
+    通过关键词搜索抖音平台的视频内容,支持多种排序和筛选方式。
+
+    Args:
+        keyword: 搜索关键词
+        content_type: 内容类型(可选:视频/图文, 默认 "视频")
+        sort_type: 排序方式(可选:综合排序/最新发布/最多点赞, 默认 "综合排序")
+        publish_time: 发布时间范围(可选:不限/一天内/一周内/半年内, 默认 "不限")
+        cursor: 分页游标,用于获取下一页结果,默认 "0"
+        account_id: 账号ID(可选)
+        timeout: 超时时间(秒),默认 60
+
+    Returns:
+        JSON 字符串,包含 output(文本摘要)和 search_results(结构化列表)。
+        search_results 中每项含 aweme_id、desc、author、statistics。
+        使用 next_cursor 可获取下一页。
+    """
+    start_time = time.time()
+    request_timeout = timeout if timeout is not None else DEFAULT_TIMEOUT
+
+    try:
+        global _last_request_monotonic
+        async with _rate_limit_lock:
+            now_mono = time.monotonic()
+            wait_seconds = _MIN_REQUEST_INTERVAL_SECONDS - (now_mono - _last_request_monotonic)
+            if wait_seconds > 0:
+                await asyncio.sleep(wait_seconds)
+            _last_request_monotonic = time.monotonic()
+
+        payload = {
+            "keyword": keyword,
+            "content_type": content_type,
+            "sort_type": sort_type,
+            "publish_time": publish_time,
+            "cursor": cursor,
+            "account_id": account_id,
+        }
+
+        async with httpx.AsyncClient(timeout=request_timeout) as client:
+            response = await client.post(
+                DOUYIN_SEARCH_API,
+                json=payload,
+                headers={"Content-Type": "application/json"},
+            )
+            response.raise_for_status()
+            data = response.json()
+
+        data_block = data.get("data", {}) if isinstance(data.get("data"), dict) else {}
+        items = data_block.get("data", []) if isinstance(data_block.get("data"), list) else []
+        has_more = bool(data_block.get("has_more", False))
+        cursor_value = str(data_block.get("next_cursor", ""))
+
+        duration_ms = int((time.time() - start_time) * 1000)
+        logger.info(
+            "douyin_search completed: keyword=%s results=%d has_more=%s duration_ms=%d",
+            keyword,
+            len(items),
+            has_more,
+            duration_ms,
+        )
+
+        return _success_result(keyword, data, items, has_more, cursor_value, duration_ms)
+
+    except httpx.HTTPStatusError as e:
+        logger.error(
+            "douyin_search HTTP error: keyword=%s status=%d",
+            keyword,
+            e.response.status_code,
+        )
+        return _error_result(f"HTTP {e.response.status_code}: {e.response.text}")
+    except httpx.TimeoutException:
+        logger.error("douyin_search timeout: keyword=%s timeout=%s", keyword, request_timeout)
+        return _error_result(f"请求超时({request_timeout}秒)")
+    except httpx.RequestError as e:
+        logger.error("douyin_search network error: keyword=%s error=%s", keyword, e)
+        return _error_result(f"网络错误: {e}")
+    except Exception as e:
+        logger.error("douyin_search unexpected error: keyword=%s error=%s", keyword, e, exc_info=True)
+        return _error_result(f"未知错误: {e}")
+
+
+async def main() -> None:
+    result_json = await douyin_search(keyword="养老政策", account_id=DOUYIN_ACCOUNT_ID)
+    result = json.loads(result_json)
+    if "error" in result:
+        print(f"搜索失败: {result['error']}")
+    else:
+        print(result["output"])
+        print(f"\n共 {result['results_count']} 条结果")
+
+
+if __name__ == "__main__":
+    asyncio.run(main())

+ 173 - 0
agents/find_agent/tools/qwen_video_analyze.py

@@ -0,0 +1,173 @@
+"""
+千问视频解析工具
+
+通过阿里云百炼平台(DashScope 兼容模式)调用 qwen3.7-plus 模型,解析视频内容。
+"""
+from __future__ import annotations
+
+import asyncio
+import json
+import logging
+import os
+import time
+from typing import Optional
+
+from dotenv import load_dotenv
+from openai import OpenAI
+
+from supply_agent.paths import find_project_root
+from supply_agent.tools import tool
+
+logger = logging.getLogger(__name__)
+
+# DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1"
+DASHSCOPE_BASE_URL = "https://llm-33b86fznnpci2exm.cn-beijing.maas.aliyuncs.com/compatible-mode/v1"
+DEFAULT_MODEL = "qwen3.7-plus"
+DEFAULT_PROMPT = "描述这段视频的内容"
+DEFAULT_FPS = 2.0
+DEFAULT_TIMEOUT = 120.0
+
+_env_loaded = False
+
+
+def _ensure_env_loaded() -> None:
+    """从项目根目录加载 .env,使 os.getenv 能读到其中的变量。"""
+    global _env_loaded
+    if _env_loaded:
+        return
+    load_dotenv(find_project_root() / ".env")
+    _env_loaded = True
+
+
+def _get_client() -> OpenAI:
+    _ensure_env_loaded()
+    api_key = os.getenv("DASHSCOPE_API_KEY")
+    if not api_key:
+        raise ValueError("未设置环境变量 DASHSCOPE_API_KEY")
+    return OpenAI(api_key=api_key, base_url=DASHSCOPE_BASE_URL)
+
+
+def _analyze_video_sync(
+    video_url: str,
+    prompt: str,
+    fps: float,
+    model: str,
+    timeout: float,
+) -> str:
+    client = _get_client()
+    completion = client.chat.completions.create(
+        model=model,
+        messages=[
+            {
+                "role": "user",
+                "content": [
+                    {
+                        "type": "video_url",
+                        "video_url": {"url": video_url},
+                        "fps": fps,
+                    },
+                    {"type": "text", "text": prompt},
+                ],
+            }
+        ],
+        timeout=timeout,
+    )
+    return completion.choices[0].message.content or ""
+
+
+def _success_result(
+    video_url: str,
+    prompt: str,
+    content: str,
+    model: str,
+    duration_ms: int,
+) -> str:
+    payload = {
+        "title": "视频解析结果",
+        "video_url": video_url,
+        "prompt": prompt,
+        "model": model,
+        "content": content,
+        "output": content,
+        "duration_ms": duration_ms,
+    }
+    return json.dumps(payload, ensure_ascii=False)
+
+
+def _error_result(error: str, *, title: str = "视频解析失败") -> str:
+    return json.dumps({"error": error, "title": title}, ensure_ascii=False)
+
+
+@tool
+async def qwen_video_analyze(
+    video_url: str,
+    prompt: str = DEFAULT_PROMPT,
+    fps: float = DEFAULT_FPS,
+    model: str = DEFAULT_MODEL,
+    timeout: Optional[float] = None,
+) -> str:
+    """
+    千问视频内容解析
+
+    通过阿里云百炼平台调用 qwen3.7-plus 模型,分析视频 URL 并返回文字描述。
+    需要设置环境变量 DASHSCOPE_API_KEY。
+
+    Args:
+        video_url: 视频地址(需公网可访问的 mp4 等格式)
+        prompt: 解析提示词,默认 "描述这段视频的内容"
+        fps: 视频抽帧频率,默认 2(每秒采样 2 帧)
+        model: 模型名称,默认 "qwen3.7-plus"
+        timeout: 请求超时时间(秒),默认 120
+
+    Returns:
+        JSON 字符串,包含 content(解析文本)和 output(同 content,供 LLM 阅读)。
+    """
+    start_time = time.time()
+    request_timeout = timeout if timeout is not None else DEFAULT_TIMEOUT
+
+    try:
+        content = await asyncio.to_thread(
+            _analyze_video_sync,
+            video_url,
+            prompt,
+            fps,
+            model,
+            request_timeout,
+        )
+
+        duration_ms = int((time.time() - start_time) * 1000)
+        logger.info(
+            "qwen_video_analyze completed: video_url=%s model=%s duration_ms=%d",
+            video_url,
+            model,
+            duration_ms,
+        )
+        return _success_result(video_url, prompt, content, model, duration_ms)
+
+    except ValueError as e:
+        logger.error("qwen_video_analyze config error: %s", e)
+        return _error_result(str(e))
+    except Exception as e:
+        logger.error(
+            "qwen_video_analyze error: video_url=%s error=%s",
+            video_url,
+            e,
+            exc_info=True,
+        )
+        return _error_result(str(e))
+
+
+async def main() -> None:
+    test_url = os.getenv("TEST_VIDEO_URL", "https://www.douyin.com/aweme/v1/play/?video_id=v0200fg10000d4lbj67og65hqvt943u0&ratio=1080p&line=0")
+    # test_url = os.getenv("TEST_VIDEO_URL", "http://rescdn.yishihui.com/longvideo/transcode/video/vpc/20260713/b2507233bc8fe040003d19b34c2c7d73.mp4")
+
+    result_json = await qwen_video_analyze(video_url=test_url)
+    result = json.loads(result_json)
+    if "error" in result:
+        print(f"解析失败: {result['error']}")
+    else:
+        print(result["content"])
+
+
+if __name__ == "__main__":
+    asyncio.run(main())

+ 8 - 0
jobs/init_db.py

@@ -0,0 +1,8 @@
+#!/usr/bin/env python3
+"""CLI entry point to initialize database tables."""
+
+from supply_infra.db import init_db
+
+if __name__ == "__main__":
+    init_db()
+    print("Database tables created.")

+ 14 - 0
jobs/run_scheduler.py

@@ -0,0 +1,14 @@
+#!/usr/bin/env python3
+"""CLI entry point for the scheduler."""
+
+import logging
+
+from supply_infra.scheduler import run_scheduler
+
+logging.basicConfig(
+    level=logging.INFO,
+    format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
+)
+
+if __name__ == "__main__":
+    run_scheduler()

+ 9 - 1
pyproject.toml

@@ -14,9 +14,14 @@ dependencies = [
     "pydantic-settings>=2.0",
     "pydantic-settings>=2.0",
     "httpx>=0.27.0",
     "httpx>=0.27.0",
     "rich>=13.0",
     "rich>=13.0",
+    "sqlalchemy>=2.0",
+    "pymysql>=1.1",
+    "apscheduler>=3.10",
+    "python-dotenv>=1.0",
 ]
 ]
 
 
 [project.optional-dependencies]
 [project.optional-dependencies]
+odps = ["pyodps>=0.12"]
 dev = [
 dev = [
     "pytest>=8.0",
     "pytest>=8.0",
     "pytest-asyncio>=0.24",
     "pytest-asyncio>=0.24",
@@ -24,7 +29,10 @@ dev = [
 ]
 ]
 
 
 [tool.hatch.build.targets.wheel]
 [tool.hatch.build.targets.wheel]
-packages = ["supply_agent"]
+packages = ["supply_agent", "supply_infra", "agents"]
+
+[project.scripts]
+supply-scheduler = "supply_infra.scheduler.app:run_scheduler"
 
 
 [tool.ruff]
 [tool.ruff]
 line-length = 100
 line-length = 100

+ 18 - 0
requirements.txt

@@ -0,0 +1,18 @@
+# SupplyAgent runtime dependencies
+
+openai>=1.50.0
+pydantic>=2.0
+pydantic-settings>=2.0
+httpx>=0.27.0
+rich>=13.0
+python-dotenv>=1.0.0
+
+# Database
+sqlalchemy>=2.0
+pymysql>=1.1
+
+# Scheduler
+apscheduler>=3.10
+
+# ODPS (optional)
+pyodps>=0.12

+ 21 - 11
supply_agent/agent/core.py

@@ -12,17 +12,27 @@ from supply_agent.types import AgentEvent, AgentEventType, AgentResult, Message,
 
 
 
 
 DEFAULT_SYSTEM_PROMPT = """\
 DEFAULT_SYSTEM_PROMPT = """\
-You are a helpful AI assistant powered by SupplyAgent.
-
-You have access to tools to help accomplish tasks. Use them when needed.
-Think step by step, call tools to gather information or take actions, \
-and provide clear final answers.
-
-## Skills (IMPORTANT)
-When the system prompt lists Available Skills, you MUST call `load_skill` \
-with the matching skill name BEFORE starting the task. \
-Skills contain specialized instructions that you must follow.
-Do NOT guess the output format — load the skill first.
+你是一个通用 AI 助手,根据用户的要求理解任务并完成目标。
+
+## 工作方式
+1. 先理解用户真正想要什么,必要时先澄清关键信息
+2. 将复杂任务拆解为可执行的步骤,逐步推进
+3. 需要外部信息或操作时,主动调用可用工具
+4. 完成后给出清晰、可直接使用的结果
+
+## 工具使用
+- 你有权访问一组工具,仅在确实需要时调用
+- 调用前想清楚:需要什么输入、期望得到什么输出
+- 工具返回后,结合结果继续推理,不要重复无效调用
+
+## 技能(Skills)
+- 若系统提示中列出了可用 Skills,且任务与某项技能匹配,先调用 `load_skill` 加载对应技能
+- 技能包含专业流程和输出规范,加载后严格遵循,不要自行猜测格式
+
+## 输出要求
+- 回答紧扣用户要求,避免无关内容
+- 结构清晰,重点突出,便于阅读和直接使用
+- 无法完成时,说明原因并给出可行的替代方案
 """
 """
 
 
 
 

+ 1 - 0
supply_infra/__init__.py

@@ -0,0 +1 @@
+"""SupplyAgent shared infrastructure: database, ODPS, scheduler."""

+ 54 - 0
supply_infra/config.py

@@ -0,0 +1,54 @@
+from __future__ import annotations
+
+from functools import lru_cache
+from urllib.parse import quote_plus
+
+from pydantic import Field
+from pydantic_settings import BaseSettings, SettingsConfigDict
+
+
+class InfraSettings(BaseSettings):
+    """Shared infrastructure settings."""
+
+    model_config = SettingsConfigDict(
+        env_file=".env",
+        env_file_encoding="utf-8",
+        extra="ignore",
+    )
+
+    # MySQL
+    mysql_host: str = Field(default="127.0.0.1", alias="MYSQL_HOST")
+    mysql_port: int = Field(default=3306, alias="MYSQL_PORT")
+    mysql_user: str = Field(default="root", alias="MYSQL_USER")
+    mysql_password: str = Field(default="", alias="MYSQL_PASSWORD")
+    mysql_database: str = Field(default="supply_agent", alias="MYSQL_DATABASE")
+    mysql_pool_size: int = Field(default=5, alias="MYSQL_POOL_SIZE")
+    mysql_echo: bool = Field(default=False, alias="MYSQL_ECHO")
+
+    # ODPS (MaxCompute)
+    odps_access_id: str = Field(default="", alias="ODPS_ACCESS_ID")
+    odps_access_key: str = Field(default="", alias="ODPS_ACCESS_KEY")
+    odps_project: str = Field(default="", alias="ODPS_PROJECT")
+    odps_endpoint: str = Field(
+        default="https://service.cn.maxcompute.aliyun.com/api",
+        alias="ODPS_ENDPOINT",
+    )
+
+    # Scheduler
+    scheduler_timezone: str = Field(default="Asia/Shanghai", alias="SCHEDULER_TIMEZONE")
+    scheduler_enabled: bool = Field(default=True, alias="SCHEDULER_ENABLED")
+
+    @property
+    def mysql_url(self) -> str:
+        user = quote_plus(self.mysql_user)
+        password = quote_plus(self.mysql_password)
+        return (
+            f"mysql+pymysql://{user}:{password}"
+            f"@{self.mysql_host}:{self.mysql_port}/{self.mysql_database}"
+            f"?charset=utf8mb4"
+        )
+
+
+@lru_cache
+def get_infra_settings() -> InfraSettings:
+    return InfraSettings()

+ 6 - 0
supply_infra/db/__init__.py

@@ -0,0 +1,6 @@
+"""Database layer."""
+
+from supply_infra.db.base import Base
+from supply_infra.db.session import get_session, init_db
+
+__all__ = ["Base", "get_session", "init_db"]

+ 26 - 0
supply_infra/db/base.py

@@ -0,0 +1,26 @@
+from __future__ import annotations
+
+from datetime import datetime
+
+from sqlalchemy import DateTime, func
+from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
+
+
+class Base(DeclarativeBase):
+    """SQLAlchemy declarative base for all ORM models."""
+
+
+class TimestampMixin:
+    """created_at / updated_at mixin for entities."""
+
+    created_at: Mapped[datetime] = mapped_column(
+        DateTime,
+        server_default=func.now(),
+        nullable=False,
+    )
+    updated_at: Mapped[datetime] = mapped_column(
+        DateTime,
+        server_default=func.now(),
+        onupdate=func.now(),
+        nullable=False,
+    )

+ 6 - 0
supply_infra/db/models/__init__.py

@@ -0,0 +1,6 @@
+"""ORM entity models — one file per table."""
+
+from supply_infra.db.models.global_tree_category import GlobalTreeCategory
+from supply_infra.db.models.global_tree_element import GlobalTreeElement
+
+__all__ = ["GlobalTreeCategory", "GlobalTreeElement"]

+ 32 - 0
supply_infra/db/models/global_tree_category.py

@@ -0,0 +1,32 @@
+from __future__ import annotations
+
+from datetime import datetime
+
+from sqlalchemy import BigInteger, Integer, String, Text, func
+from sqlalchemy.orm import Mapped, mapped_column
+
+from supply_infra.db.base import Base
+
+
+class GlobalTreeCategory(Base):
+    """全局树分类 — 从 ODPS public_pattern_mining_category 同步。"""
+
+    __tablename__ = "global_tree_category"
+
+    id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
+    name: Mapped[str | None] = mapped_column(String(128), nullable=True, comment="名称")
+    description: Mapped[str | None] = mapped_column(Text, nullable=True, comment="描述")
+    level: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="层级")
+    parent_id: Mapped[int | None] = mapped_column(BigInteger, nullable=True, comment="父节点")
+    is_delete: Mapped[int] = mapped_column(Integer, default=0, nullable=False, comment="是否删除0-正常 1-删除")
+    create_time: Mapped[datetime | None] = mapped_column(
+        nullable=True,
+        server_default=func.now(),
+        comment="创建时间",
+    )
+    update_time: Mapped[datetime] = mapped_column(
+        nullable=False,
+        server_default=func.now(),
+        onupdate=func.now(),
+        comment="更新时间",
+    )

+ 34 - 0
supply_infra/db/models/global_tree_element.py

@@ -0,0 +1,34 @@
+from __future__ import annotations
+
+from datetime import datetime
+
+from sqlalchemy import BigInteger, Index, Integer, String, UniqueConstraint, func
+from sqlalchemy.orm import Mapped, mapped_column
+
+from supply_infra.db.base import Base
+
+
+class GlobalTreeElement(Base):
+    """全局树元素 — 从 ODPS pattern_mining_element 同步。"""
+
+    __tablename__ = "global_tree_element"
+    __table_args__ = (
+        UniqueConstraint("name", "category_id", name="global_tree_element_pk"),
+        Index("idx_category_id", "category_id"),
+    )
+
+    id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
+    name: Mapped[str | None] = mapped_column(String(128), nullable=True, comment="名称")
+    category_id: Mapped[int] = mapped_column(BigInteger, nullable=False, comment="分类id")
+    is_delete: Mapped[int] = mapped_column(Integer, default=0, nullable=False, comment="是否删除 0-正常 1-删除")
+    create_time: Mapped[datetime] = mapped_column(
+        nullable=False,
+        server_default=func.now(),
+        comment="创建时间",
+    )
+    update_time: Mapped[datetime] = mapped_column(
+        nullable=False,
+        server_default=func.now(),
+        onupdate=func.now(),
+        comment="更新时间",
+    )

+ 11 - 0
supply_infra/db/repositories/__init__.py

@@ -0,0 +1,11 @@
+"""Data access layer (Repository pattern)."""
+
+from supply_infra.db.repositories.base import BaseRepository
+from supply_infra.db.repositories.global_tree_category_repo import GlobalTreeCategoryRepository
+from supply_infra.db.repositories.global_tree_element_repo import GlobalTreeElementRepository
+
+__all__ = [
+    "BaseRepository",
+    "GlobalTreeCategoryRepository",
+    "GlobalTreeElementRepository",
+]

+ 34 - 0
supply_infra/db/repositories/base.py

@@ -0,0 +1,34 @@
+from __future__ import annotations
+
+from typing import Generic, TypeVar
+
+from sqlalchemy import select
+from sqlalchemy.orm import Session
+
+from supply_infra.db.base import Base
+
+T = TypeVar("T", bound=Base)
+
+
+class BaseRepository(Generic[T]):
+    """Generic CRUD repository — subclass per entity."""
+
+    model: type[T]
+
+    def __init__(self, session: Session) -> None:
+        self.session = session
+
+    def get_by_id(self, id: int) -> T | None:
+        return self.session.get(self.model, id)
+
+    def list_all(self, *, limit: int = 100, offset: int = 0) -> list[T]:
+        stmt = select(self.model).limit(limit).offset(offset)
+        return list(self.session.scalars(stmt).all())
+
+    def add(self, entity: T) -> T:
+        self.session.add(entity)
+        self.session.flush()
+        return entity
+
+    def delete(self, entity: T) -> None:
+        self.session.delete(entity)

+ 27 - 0
supply_infra/db/repositories/global_tree_category_repo.py

@@ -0,0 +1,27 @@
+from __future__ import annotations
+
+from sqlalchemy.dialects.mysql import insert
+
+from supply_infra.db.models.global_tree_category import GlobalTreeCategory
+from supply_infra.db.repositories.base import BaseRepository
+
+_BATCH_SIZE = 1000
+
+
+class GlobalTreeCategoryRepository(BaseRepository[GlobalTreeCategory]):
+    """全局树分类表 — INSERT IGNORE 按 id 主键去重。"""
+
+    model = GlobalTreeCategory
+
+    def bulk_insert_ignore(self, rows: list[dict]) -> int:
+        """批量插入,MySQL 自动忽略 id 重复行。"""
+        if not rows:
+            return 0
+
+        inserted = 0
+        for i in range(0, len(rows), _BATCH_SIZE):
+            batch = rows[i : i + _BATCH_SIZE]
+            stmt = insert(GlobalTreeCategory).values(batch).prefix_with("IGNORE")
+            result = self.session.execute(stmt)
+            inserted += result.rowcount
+        return inserted

+ 27 - 0
supply_infra/db/repositories/global_tree_element_repo.py

@@ -0,0 +1,27 @@
+from __future__ import annotations
+
+from sqlalchemy.dialects.mysql import insert
+
+from supply_infra.db.models.global_tree_element import GlobalTreeElement
+from supply_infra.db.repositories.base import BaseRepository
+
+_BATCH_SIZE = 1000
+
+
+class GlobalTreeElementRepository(BaseRepository[GlobalTreeElement]):
+    """全局树元素表 — INSERT IGNORE 按 (name, category_id) 唯一键去重。"""
+
+    model = GlobalTreeElement
+
+    def bulk_insert_ignore(self, rows: list[dict]) -> int:
+        """批量插入,MySQL 自动忽略 (name, category_id) 重复行。"""
+        if not rows:
+            return 0
+
+        inserted = 0
+        for i in range(0, len(rows), _BATCH_SIZE):
+            batch = rows[i : i + _BATCH_SIZE]
+            stmt = insert(GlobalTreeElement).values(batch).prefix_with("IGNORE")
+            result = self.session.execute(stmt)
+            inserted += result.rowcount
+        return inserted

+ 52 - 0
supply_infra/db/session.py

@@ -0,0 +1,52 @@
+from __future__ import annotations
+
+from collections.abc import Generator
+from contextlib import contextmanager
+from typing import Any
+
+from sqlalchemy import create_engine
+from sqlalchemy.orm import Session, sessionmaker
+
+from supply_infra.config import get_infra_settings
+from supply_infra.db.base import Base
+
+_engine: Any = None
+_SessionLocal: sessionmaker[Session] | None = None
+
+
+def get_engine():
+    """Lazy-init SQLAlchemy engine (singleton)."""
+    global _engine, _SessionLocal
+    if _engine is None:
+        settings = get_infra_settings()
+        _engine = create_engine(
+            settings.mysql_url,
+            pool_size=settings.mysql_pool_size,
+            pool_pre_ping=True,
+            echo=settings.mysql_echo,
+        )
+        _SessionLocal = sessionmaker(bind=_engine, autoflush=False, autocommit=False)
+    return _engine
+
+
+def init_db() -> None:
+    """Create all tables (dev / first-run). Import models before calling."""
+    import supply_infra.db.models  # noqa: F401 — register all models
+
+    Base.metadata.create_all(bind=get_engine())
+
+
+@contextmanager
+def get_session() -> Generator[Session, None, None]:
+    """Provide a transactional database session."""
+    get_engine()
+    assert _SessionLocal is not None
+    session = _SessionLocal()
+    try:
+        yield session
+        session.commit()
+    except Exception:
+        session.rollback()
+        raise
+    finally:
+        session.close()

+ 5 - 0
supply_infra/odps/__init__.py

@@ -0,0 +1,5 @@
+"""ODPS (MaxCompute) client."""
+
+from supply_infra.odps.client import ODPSClient, get_odps_client
+
+__all__ = ["ODPSClient", "get_odps_client"]

+ 96 - 0
supply_infra/odps/client.py

@@ -0,0 +1,96 @@
+from __future__ import annotations
+
+import logging
+from functools import lru_cache
+from typing import Any
+
+from supply_infra.config import get_infra_settings
+
+logger = logging.getLogger(__name__)
+
+
+class ODPSClient:
+    """ODPS / MaxCompute 查询客户端封装。"""
+
+    def __init__(
+        self,
+        access_id: str,
+        access_key: str,
+        project: str,
+        endpoint: str,
+    ) -> None:
+        self.access_id = access_id
+        self.access_key = access_key
+        self.project = project
+        self.endpoint = endpoint
+        self._client: Any = None
+
+    def _get_client(self) -> Any:
+        if self._client is None:
+            try:
+                from odps import ODPS
+            except ImportError as e:
+                raise ImportError(
+                    "pyodps 未安装,请执行: pip install pyodps"
+                ) from e
+            self._client = ODPS(
+                self.access_id,
+                self.access_key,
+                project=self.project,
+                endpoint=self.endpoint,
+            )
+        return self._client
+
+    def execute_sql(self, sql: str) -> list[dict[str, Any]]:
+        """执行 SQL 并返回字典列表。"""
+        client = self._get_client()
+        logger.info("ODPS executing SQL: %s", sql[:200])
+        instance = client.execute_sql(sql)
+        with instance.open_reader() as reader:
+            columns = [col.name for col in reader._schema.columns]  # type: ignore[attr-defined]
+            return [dict(zip(columns, row.values)) for row in reader]
+        return []
+
+
+    def fetch_pattern_mining_elements(self, bizdate: str) -> list[dict[str, Any]]:
+        """拉取 pattern_mining_element 元素(name, category_id)。"""
+        sql = f"""
+        SELECT  t1.name
+                ,t1.category_id
+        FROM    loghubods.pattern_mining_element t1
+        LEFT JOIN loghubods.post t2
+        ON      t1.post_id = t2.post_id
+        AND     t2.dt = '{bizdate}'
+        WHERE   t1.dt = '{bizdate}'
+        AND     t1.execution_id = 401
+        AND     t1.source_table = 'post_decode_topic_point_element'
+        AND     t1.name is NOT NULL
+        AND     t1.category_id IS NOT NULL
+        AND     t1.element_type = '实质'
+        AND     t2.platform = 'piaoquan'
+        GROUP BY t1.name
+                 ,t1.category_id
+        """
+        return self.execute_sql(sql)
+
+    def fetch_pattern_mining_categories(self, bizdate: str) -> list[dict[str, Any]]:
+        """拉取 public_pattern_mining_category 分类。"""
+        sql = f"""
+        SELECT  id,name,description,level,parent_id
+        FROM    loghubods.public_pattern_mining_category
+        WHERE   dt = '{bizdate}'
+        AND     execution_id = 401
+        AND     source_type = '实质'
+        """
+        return self.execute_sql(sql)
+
+
+@lru_cache
+def get_odps_client() -> ODPSClient:
+    settings = get_infra_settings()
+    return ODPSClient(
+        access_id=settings.odps_access_id,
+        access_key=settings.odps_access_key,
+        project=settings.odps_project,
+        endpoint=settings.odps_endpoint,
+    )

+ 5 - 0
supply_infra/scheduler/__init__.py

@@ -0,0 +1,5 @@
+"""Scheduler for periodic jobs."""
+
+from supply_infra.scheduler.app import create_scheduler, run_scheduler
+
+__all__ = ["create_scheduler", "run_scheduler"]

+ 47 - 0
supply_infra/scheduler/app.py

@@ -0,0 +1,47 @@
+from __future__ import annotations
+
+import logging
+
+from apscheduler.schedulers.blocking import BlockingScheduler
+from apscheduler.triggers.cron import CronTrigger
+
+from supply_infra.config import get_infra_settings
+from supply_infra.scheduler.jobs.sync_global_tree_odps_to_mysql import sync_global_tree_odps_to_mysql
+
+logger = logging.getLogger(__name__)
+
+
+def create_scheduler() -> BlockingScheduler:
+    """Create and configure the scheduler with all registered jobs."""
+    settings = get_infra_settings()
+    scheduler = BlockingScheduler(timezone=settings.scheduler_timezone)
+
+    # 每天凌晨 2:30 从 ODPS 同步全局树元素与分类到 MySQL
+    scheduler.add_job(
+        sync_global_tree_odps_to_mysql,
+        trigger=CronTrigger(hour=2, minute=30),
+        id="sync_global_tree_odps_to_mysql",
+        name="ODPS → MySQL 全局树同步",
+        replace_existing=True,
+    )
+
+    logger.info("Scheduler configured with %d job(s)", len(scheduler.get_jobs()))
+    return scheduler
+
+
+def run_scheduler() -> None:
+    """Start the blocking scheduler (CLI entry point)."""
+    settings = get_infra_settings()
+    if not settings.scheduler_enabled:
+        logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
+        return
+
+    scheduler = create_scheduler()
+    logger.info("Starting scheduler...")
+    for job in scheduler.get_jobs():
+        logger.info("  - %s | next run: %s", job.name, job.next_run_time)
+
+    try:
+        scheduler.start()
+    except (KeyboardInterrupt, SystemExit):
+        logger.info("Scheduler stopped.")

+ 1 - 0
supply_infra/scheduler/jobs/__init__.py

@@ -0,0 +1 @@
+"""Scheduled job definitions."""

+ 122 - 0
supply_infra/scheduler/jobs/sync_global_tree_odps_to_mysql.py

@@ -0,0 +1,122 @@
+"""
+定时任务:从 ODPS 同步全局树元素与分类到 MySQL。
+
+流程:
+1. 拉取 pattern_mining_element(name, category_id)
+2. 拉取 public_pattern_mining_category 全量分类
+3. 筛选元素关联的分类,并沿 parent_id 向上追溯至顶层
+4. 分别批量 INSERT IGNORE 写入 global_tree_element / global_tree_category
+"""
+from __future__ import annotations
+
+import logging
+from datetime import datetime, timedelta
+from typing import Any
+
+from supply_infra.db.repositories.global_tree_category_repo import GlobalTreeCategoryRepository
+from supply_infra.db.repositories.global_tree_element_repo import GlobalTreeElementRepository
+from supply_infra.db.session import get_session
+from supply_infra.odps.client import get_odps_client
+
+logger = logging.getLogger(__name__)
+
+
+def _collect_categories_with_ancestors(
+    seed_ids: set[int],
+    category_by_id: dict[int, dict[str, Any]],
+) -> list[dict[str, Any]]:
+    """从种子 category_id 出发,沿 parent_id 向上追溯,收集完整祖先链。"""
+    collected: dict[int, dict[str, Any]] = {}
+    pending = set(seed_ids)
+
+    while pending:
+        cat_id = pending.pop()
+        if cat_id in collected:
+            continue
+
+        category = category_by_id.get(cat_id)
+        if category is None:
+            continue
+
+        collected[cat_id] = category
+        parent_id = category.get("parent_id")
+        if parent_id is not None and int(parent_id) != 0 and int(parent_id) not in collected:
+            pending.add(int(parent_id))
+
+    return list(collected.values())
+
+
+def _to_element_rows(elements: list[dict[str, Any]]) -> list[dict[str, Any]]:
+    return [
+        {"name": str(row["name"]), "category_id": int(row["category_id"])}
+        for row in elements
+        if row.get("name") is not None and row.get("category_id") is not None
+    ]
+
+
+def _to_category_rows(categories: list[dict[str, Any]]) -> list[dict[str, Any]]:
+    rows: list[dict[str, Any]] = []
+    for cat in categories:
+        parent_id = cat.get("parent_id")
+        rows.append(
+            {
+                "id": int(cat["id"]),
+                "name": str(cat["name"]) if cat.get("name") is not None else None,
+                "description": cat.get("description"),
+                "level": int(cat["level"]) if cat.get("level") is not None else None,
+                "parent_id": int(parent_id) if parent_id is not None else None,
+            }
+        )
+    return rows
+
+
+def sync_global_tree_odps_to_mysql(partition_date: str | None = None) -> dict:
+    """
+    从 ODPS 拉取全局树元素与分类,写入 MySQL。
+
+    Args:
+        partition_date: 分区日期 (YYYYMMDD),默认昨天
+    """
+    if partition_date is None:
+        partition_date = (datetime.now() - timedelta(days=1)).strftime("%Y%m%d")
+
+    logger.info("Starting global tree ODPS → MySQL sync for partition: %s", partition_date)
+
+    odps = get_odps_client()
+    raw_elements = odps.fetch_pattern_mining_elements(partition_date)
+    raw_categories = odps.fetch_pattern_mining_categories(partition_date)
+    logger.info(
+        "Fetched %d elements and %d categories from ODPS",
+        len(raw_elements),
+        len(raw_categories),
+    )
+
+    elements = _to_element_rows(raw_elements)
+    category_by_id = {int(c["id"]): c for c in raw_categories if c.get("id") is not None}
+
+    seed_ids = {row["category_id"] for row in elements}
+    categories = _collect_categories_with_ancestors(seed_ids, category_by_id)
+    category_rows = _to_category_rows(categories)
+
+    logger.info(
+        "Resolved %d elements and %d categories (with ancestors) for insert",
+        len(elements),
+        len(category_rows),
+    )
+
+    with get_session() as session:
+        element_repo = GlobalTreeElementRepository(session)
+        category_repo = GlobalTreeCategoryRepository(session)
+        element_inserted = element_repo.bulk_insert_ignore(elements)
+        category_inserted = category_repo.bulk_insert_ignore(category_rows)
+
+    result = {
+        "partition_date": partition_date,
+        "elements_fetched": len(elements),
+        "elements_inserted": element_inserted,
+        "categories_fetched": len(category_rows),
+        "categories_inserted": category_inserted,
+        "synced_at": datetime.now().isoformat(),
+    }
+    logger.info("Global tree sync completed: %s", result)
+    return result