多 Agent 数据管道:从采集到报告的自动化流水线
采集 → 分析 → 报告自动化 | Hermes Cron 定时触发 → Codex 分析数据 → Claude Code 生成报告
简介
在现代数据驱动的工作流中,数据处理往往需要多个环节的协作:从不同来源采集数据、清洗转换、统计分析、生成报告。传统方式下,这些步骤通常需要不同的人员或工具来完成,耗时且容易出错。
如果我们能将这些步骤交给 AI Agent 自动完成,让数据采集 Agent 负责抓取、数据分析 Agent 负责洞察、报告生成 Agent 负责呈现,并且通过 Hermes Cron 定时触发整个流程——这将彻底改变我们处理数据的方式。
本文将详细介绍如何构建这样一个多 Agent 数据管道系统。
一、管道架构概览
1.1 整体架构
┌──────────────────────────────────────────────────────────────┐
│ Hermes Cron 调度器 │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ 每日采集 │ │ 周报分析 │ │ 月度报告 │ │ 异常触发 │ │
│ └────┬────┘ └────┬────┘ └────┬────┘ └────┬────┘ │
└───────┼────────────┼────────────┼────────────┼──────────────┘
│ │ │ │
▼ ▼ ▼ ▼
┌──────────────────────────────────────────────────────────────┐
│ 数据管道 (Pipeline) │
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │
│ │ 采集Agent │───▶│ 清洗Agent │───▶│ 分析Agent │───▶│报告Agent│ │
│ │ (Codex) │ │ (Codex) │ │ (Codex) │ │(Claude) │ │
│ └──────────┘ └──────────┘ └──────────┘ └────────┘ │
│ │ │ │ │ │
│ ▼ ▼ ▼ ▼ │
│ 原始数据 清洗数据 分析结果 PDF/HTML │
└──────────────────────────────────────────────────────────────┘
│ │ │ │
▼ ▼ ▼ ▼
┌──────────┐ ┌──────────────┐ ┌────────────┐ ┌────────────┐
│ API/爬虫 │ │ 数据湖/S3 │ │ 结果数据库 │ │ 邮件/微信 │
└──────────┘ └──────────────┘ └────────────┘ └────────────┘1.2 管道阶段定义
| 阶段 | 负责 Agent | 核心能力 | 输出 |
|---|---|---|---|
| 采集 | Codex Agent | API 调用、网页爬取、数据库查询 | 原始数据 JSON/CSV |
| 清洗 | Codex Agent | 数据清洗、格式转换、异常值处理 | 清洗后数据集 |
| 分析 | Codex Agent | 统计分析、趋势检测、异常识别 | 分析结果 + 洞察 |
| 报告 | Claude Code | 报告编写、图表生成、排版美化 | PDF/HTML/Markdown |
二、Hermes Cron 定时调度
2.1 Cron 配置
Hermes Cron 提供了强大的定时任务管理能力,支持 cron 表达式、间隔触发、事件驱动等多种模式:
# config/cron-pipelines.yaml
cron:
timezone: "Asia/Shanghai"
jobs:
# 每日数据采集 - 早上 8 点
- name: "daily-data-collection"
schedule: "0 8 * * *"
enabled: true
pipeline: "data-collection"
config:
sources:
- type: "api"
url: "https://api.example.com/metrics"
auth: "${API_KEY}"
- type: "database"
connection: "${DB_DSN}"
query: "SELECT * FROM daily_stats WHERE date = CURRENT_DATE - 1"
- type: "web"
url: "https://competitor.example.com/pricing"
scraper: "pricing-scraper"
# 每周数据分析 - 周一早上 9 点
- name: "weekly-data-analysis"
schedule: "0 9 * * 1"
enabled: true
pipeline: "data-analysis"
depends_on: "daily-data-collection"
config:
analysis_types:
- "trend_analysis"
- "anomaly_detection"
- "correlation_analysis"
lookback_days: 7
# 月度报告生成 - 每月 1 号 10 点
- name: "monthly-report"
schedule: "0 10 1 * *"
enabled: true
pipeline: "report-generation"
depends_on: "weekly-data-analysis"
config:
template: "monthly-report-v2"
formats: ["pdf", "html", "markdown"]
recipients:
- type: "email"
addresses: ["team@example.com"]
- type: "wechat"
group: "技术团队群"
# 异常数据实时触发
- name: "anomaly-alert"
trigger: "event"
event_pattern: "data.anomaly.detected"
enabled: true
pipeline: "anomaly-response"
config:
alert_channels: ["wechat", "slack", "sms"]
severity_threshold: "high"2.2 Pipeline 定义
# config/pipelines/data-collection.yaml
pipeline:
name: "data-collection"
version: "1.0"
stages:
- name: "collect"
agent: "codex-collector"
timeout: 600s
retry: 3
- name: "validate"
agent: "codex-validator"
timeout: 120s
retry: 2
- name: "store"
agent: "storage-handler"
timeout: 60s
retry: 3
# 阶段间数据传递
data_flow:
collect.output -> validate.input
validate.output -> store.input
# 错误处理
on_failure:
notify: ["wechat:admin", "slack:#data-alerts"]
fallback: "use_last_valid_data"三、Codex Agent 数据分析
3.1 数据采集 Agent
Codex 擅长理解和编写代码,非常适合用于数据采集任务:
# src/agents/codex_collector.py
import asyncio
import aiohttp
import pandas as pd
from typing import Dict, Any, List
from dataclasses import dataclass
@dataclass
class DataSource:
name: str
source_type: str # "api", "database", "web", "file"
config: Dict[str, Any]
class CodexCollectorAgent:
"""Codex 驱动的数据采集 Agent"""
def __init__(self, config: Dict[str, Any]):
self.config = config
self.session = None
async def initialize(self):
"""初始化 HTTP 会话"""
self.session = aiohttp.ClientSession(
timeout=aiohttp.ClientTimeout(total=30)
)
async def collect(self, sources: List[DataSource]) -> Dict[str, pd.DataFrame]:
"""并发采集多个数据源"""
tasks = []
for source in sources:
if source.source_type == "api":
tasks.append(self._collect_api(source))
elif source.source_type == "database":
tasks.append(self._collect_database(source))
elif source.source_type == "web":
tasks.append(self._collect_web(source))
elif source.source_type == "file":
tasks.append(self._collect_file(source))
results = await asyncio.gather(*tasks, return_exceptions=True)
# 合并结果
collected_data = {}
for source, result in zip(sources, results):
if isinstance(result, Exception):
print(f"⚠️ {source.name} 采集失败: {result}")
else:
collected_data[source.name] = result
return collected_data
async def _collect_api(self, source: DataSource) -> pd.DataFrame:
"""API 数据采集"""
url = source.config["url"]
headers = source.config.get("headers", {})
params = source.config.get("params", {})
async with self.session.get(url, headers=headers, params=params) as resp:
resp.raise_for_status()
data = await resp.json()
# Codex 生成的解析逻辑
return self._parse_api_response(data, source.config.get("parse_rules", {}))
async def _collect_database(self, source: DataSource) -> pd.DataFrame:
"""数据库数据采集"""
import asyncpg
conn = await asyncpg.connect(source.config["dsn"])
try:
records = await conn.fetch(source.config["query"])
return pd.DataFrame([dict(r) for r in records])
finally:
await conn.close()
async def _collect_web(self, source: DataSource) -> pd.DataFrame:
"""网页数据采集"""
from bs4 import BeautifulSoup
async with self.session.get(source.config["url"]) as resp:
html = await resp.text()
soup = BeautifulSoup(html, "html.parser")
# Codex 根据页面结构生成解析代码
return self._parse_web_page(soup, source.config.get("selectors", {}))
def _parse_api_response(self, data: dict, rules: dict) -> pd.DataFrame:
"""解析 API 响应数据"""
# 支持多种响应格式
if rules.get("data_path"):
for key in rules["data_path"].split("."):
data = data[key]
if isinstance(data, list):
return pd.DataFrame(data)
elif isinstance(data, dict):
return pd.DataFrame([data])
else:
return pd.DataFrame({"value": [data]})3.2 数据分析 Agent
数据分析是管道的核心环节,Codex 可以编写复杂的分析代码:
# src/agents/codex_analyst.py
import pandas as pd
import numpy as np
from typing import Dict, Any, List
from dataclasses import dataclass
@dataclass
class AnalysisResult:
summary: Dict[str, Any]
trends: Dict[str, Any]
anomalies: List[Dict[str, Any]]
correlations: Dict[str, float]
insights: List[str]
class CodexAnalystAgent:
"""Codex 驱动的数据分析 Agent"""
def __init__(self, config: Dict[str, Any]):
self.config = config
async def analyze(self, data: Dict[str, pd.DataFrame],
analysis_types: List[str]) -> AnalysisResult:
"""执行多维度数据分析"""
results = {
"summary": {},
"trends": {},
"anomalies": [],
"correlations": {},
"insights": [],
}
for name, df in data.items():
# 1. 描述性统计
results["summary"][name] = self._descriptive_stats(df)
# 2. 趋势分析
if "trend_analysis" in analysis_types:
results["trends"][name] = self._trend_analysis(df)
# 3. 异常检测
if "anomaly_detection" in analysis_types:
anomalies = self._detect_anomalies(df)
results["anomalies"].extend(anomalies)
# 4. 相关性分析
if "correlation_analysis" in analysis_types:
results["correlations"][name] = self._correlation_analysis(df)
# 5. 生成洞察
results["insights"] = self._generate_insights(results)
return AnalysisResult(**results)
def _descriptive_stats(self, df: pd.DataFrame) -> Dict[str, Any]:
"""描述性统计"""
numeric_cols = df.select_dtypes(include=[np.number]).columns
stats = {}
for col in numeric_cols:
stats[col] = {
"count": int(df[col].count()),
"mean": round(float(df[col].mean()), 4),
"std": round(float(df[col].std()), 4),
"min": round(float(df[col].min()), 4),
"max": round(float(df[col].max()), 4),
"median": round(float(df[col].median()), 4),
"q25": round(float(df[col].quantile(0.25)), 4),
"q75": round(float(df[col].quantile(0.75)), 4),
}
return stats
def _trend_analysis(self, df: pd.DataFrame) -> Dict[str, Any]:
"""趋势分析"""
# 检测时间序列趋势
date_cols = df.select_dtypes(include=["datetime"]).columns
trends = {}
if len(date_cols) > 0:
date_col = date_cols[0]
df_sorted = df.sort_values(date_col)
numeric_cols = df.select_dtypes(include=[np.number]).columns
for col in numeric_cols:
values = df_sorted[col].values
# 简单线性回归
x = np.arange(len(values))
slope = np.polyfit(x, values, 1)[0]
trend_direction = "上升" if slope > 0 else "下降"
trends[col] = {
"direction": trend_direction,
"slope": round(float(slope), 6),
"change_pct": round(
float((values[-1] - values[0]) / values[0] * 100), 2
) if values[0] != 0 else 0,
}
return trends
def _detect_anomalies(self, df: pd.DataFrame,
method: str = "iqr") -> List[Dict[str, Any]]:
"""异常值检测"""
anomalies = []
numeric_cols = df.select_dtypes(include=[np.number]).columns
for col in numeric_cols:
values = df[col]
if method == "iqr":
Q1 = values.quantile(0.25)
Q3 = values.quantile(0.75)
IQR = Q3 - Q1
lower = Q1 - 1.5 * IQR
upper = Q3 + 1.5 * IQR
anomaly_mask = (values < lower) | (values > upper)
elif method == "zscore":
z_scores = np.abs((values - values.mean()) / values.std())
anomaly_mask = z_scores > 3
else:
continue
anomaly_indices = df[anomaly_mask].index.tolist()
for idx in anomaly_indices:
anomalies.append({
"column": col,
"value": float(df.loc[idx, col]),
"row_index": int(idx),
"method": method,
})
return anomalies
def _correlation_analysis(self, df: pd.DataFrame) -> Dict[str, float]:
"""相关性分析"""
numeric_cols = df.select_dtypes(include=[np.number]).columns
if len(numeric_cols) < 2:
return {}
corr_matrix = df[numeric_cols].corr()
# 找出强相关对
strong_correlations = {}
for i, col1 in enumerate(numeric_cols):
for j, col2 in enumerate(numeric_cols):
if i < j:
corr = corr_matrix.loc[col1, col2]
if abs(corr) > 0.7:
strong_correlations[f"{col1}↔{col2}"] = round(float(corr), 4)
return strong_correlations
def _generate_insights(self, results: Dict[str, Any]) -> List[str]:
"""基于分析结果生成自然语言洞察"""
insights = []
# 趋势洞察
for name, trends in results["trends"].items():
for col, trend in trends.items():
if abs(trend.get("change_pct", 0)) > 10:
direction = trend["direction"]
change = trend["change_pct"]
insights.append(
f"📈 {name}.{col} 呈现{direction}趋势,"
f"变化幅度 {change}%"
)
# 异常洞察
if results["anomalies"]:
anomaly_cols = set(a["column"] for a in results["anomalies"])
insights.append(
f"⚠️ 检测到 {len(results['anomalies'])} 个异常值,"
f"涉及字段: {', '.join(anomaly_cols)}"
)
# 相关性洞察
for name, corrs in results["correlations"].items():
for pair, corr in corrs.items():
strength = "强" if abs(corr) > 0.8 else "较强"
insights.append(
f"🔗 {name}.{pair} 存在{strength}相关性 (r={corr})"
)
return insights四、Claude Code 报告生成
4.1 报告模板引擎
Claude Code 负责将分析结果转化为结构化的报告:
# src/agents/claude_reporter.py
import asyncio
from pathlib import Path
from typing import Dict, Any, List
import json
class ClaudeReportAgent:
"""Claude Code 驱动的报告生成 Agent"""
def __init__(self, config: Dict[str, Any]):
self.config = config
self.template_dir = Path(config.get("template_dir", "templates/reports"))
async def generate_report(self, analysis_result,
report_config: Dict[str, Any]) -> Dict[str, str]:
"""生成多格式报告"""
# 1. 准备数据
report_data = self._prepare_report_data(analysis_result, report_config)
# 2. 加载模板
template = self._load_template(report_config.get("template", "default"))
# 3. 构建 Claude Code 提示
prompt = self._build_report_prompt(report_data, template, report_config)
# 4. 调用 Claude Code 生成报告
outputs = {}
for fmt in report_config.get("formats", ["markdown"]):
output = await self._generate_format(prompt, fmt)
outputs[fmt] = output
# 5. 保存报告
saved_paths = await self._save_reports(outputs, report_config)
return saved_paths
def _prepare_report_data(self, analysis_result, config) -> Dict:
"""准备报告数据"""
return {
"report_title": config.get("title", "数据分析报告"),
"report_date": self._get_report_date(config),
"summary": analysis_result.summary,
"trends": analysis_result.trends,
"anomalies": analysis_result.anomalies,
"correlations": analysis_result.correlations,
"insights": analysis_result.insights,
}
def _build_report_prompt(self, data: Dict, template: str,
config: Dict) -> str:
"""构建报告生成提示"""
prompt = f"""请根据以下数据分析结果生成一份专业的报告。
## 报告要求
- 标题: {data['report_title']}
- 日期: {data['report_date']}
- 格式: Markdown
- 语言: 中文
- 包含数据可视化建议
## 数据概览
{json.dumps(data['summary'], indent=2, ensure_ascii=False)}
## 趋势分析
{json.dumps(data['trends'], indent=2, ensure_ascii=False)}
## 异常检测
发现 {len(data['anomalies'])} 个异常值
## 相关性分析
{json.dumps(data['correlations'], indent=2, ensure_ascii=False)}
## 关键洞察
{chr(10).join(f"- {insight}" for insight in data['insights'])}
## 模板参考
{template}
请生成完整的报告内容,包含:
1. 执行摘要(200字以内)
2. 数据概览表格
3. 趋势分析图表建议
4. 异常值详细说明
5. 相关性解读
6. 行动建议
7. 附录(原始数据摘要)
开始生成报告。"""
return prompt
async def _generate_format(self, prompt: str, format: str) -> str:
"""调用 Claude Code 生成指定格式的报告"""
cmd = ["claude", "-p", prompt]
if format == "markdown":
cmd.extend(["--outputFormat", "text"])
elif format == "html":
prompt += "\n\n请将报告转换为 HTML 格式,包含内联 CSS 样式。"
elif format == "pdf":
prompt += "\n\n请生成适合打印的 PDF 内容。"
proc = await asyncio.create_subprocess_exec(
*cmd,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
stdout, stderr = await proc.communicate()
return stdout.decode("utf-8")4.2 报告分发
# src/distribution.py
import asyncio
import smtplib
from email.mime.multipart import MIMEMultipart
from email.mime.text import MIMEText
from email.mime.application import MIMEApplication
from pathlib import Path
from typing import List
class ReportDistributor:
"""报告分发器"""
def __init__(self, config: Dict):
self.config = config
async def distribute(self, report_paths: Dict[str, str],
recipients: List[Dict]) -> Dict[str, bool]:
"""分发报告到多个渠道"""
results = {}
for recipient in recipients:
channel = recipient["type"]
if channel == "email":
results[f"email:{recipient['addresses']}"] = \
await self._send_email(recipient, report_paths)
elif channel == "wechat":
results[f"wechat:{recipient['group']}"] = \
await self._send_wechat(recipient, report_paths)
elif channel == "slack":
results[f"slack:{recipient['channel']}"] = \
await self._send_slack(recipient, report_paths)
return results
async def _send_email(self, recipient: Dict,
report_paths: Dict[str, str]) -> bool:
"""邮件发送报告"""
try:
msg = MIMEMultipart()
msg["Subject"] = f"📊 {report_paths.get('title', '数据分析报告')}"
msg["To"] = ", ".join(recipient["addresses"])
# 附加 HTML 报告
if "html" in report_paths:
with open(report_paths["html"]) as f:
html_content = f.read()
msg.attach(MIMEText(html_content, "html"))
# 附加 PDF
if "pdf" in report_paths:
with open(report_paths["pdf"], "rb") as f:
pdf_attachment = MIMEApplication(f.read())
pdf_attachment.add_header(
"Content-Disposition", "attachment",
filename="report.pdf"
)
msg.attach(pdf_attachment)
# 发送邮件
with smtplib.SMTP(self.config["smtp_server"]) as server:
server.login(
self.config["smtp_user"],
self.config["smtp_password"]
)
server.send_message(msg)
return True
except Exception as e:
print(f"邮件发送失败: {e}")
return False五、完整管道编排
5.1 Pipeline 编排器
# src/pipeline/orchestrator.py
import asyncio
from typing import Dict, Any
from datetime import datetime
class PipelineOrchestrator:
"""管道编排器 - 串联多个 Agent"""
def __init__(self, agents: Dict[str, Any]):
self.agents = agents
self.pipeline_state = {}
async def run_pipeline(self, pipeline_name: str,
params: Dict[str, Any] = None) -> Dict:
"""运行完整的数据管道"""
start_time = datetime.now()
self.pipeline_state = {
"pipeline": pipeline_name,
"start_time": start_time.isoformat(),
"stages": {},
}
try:
# Stage 1: 数据采集
await self._run_stage("collect", params)
# Stage 2: 数据清洗
await self._run_stage("clean", params)
# Stage 3: 数据分析
await self._run_stage("analyze", params)
# Stage 4: 报告生成
await self._run_stage("report", params)
# Stage 5: 报告分发
await self._run_stage("distribute", params)
self.pipeline_state["status"] = "success"
self.pipeline_state["end_time"] = datetime.now().isoformat()
except Exception as e:
self.pipeline_state["status"] = "failed"
self.pipeline_state["error"] = str(e)
await self._handle_failure(e)
return self.pipeline_state
async def _run_stage(self, stage_name: str, params: Dict) -> None:
"""运行单个管道阶段"""
stage_start = datetime.now()
if stage_name == "collect":
collector = self.agents["codex-collector"]
await collector.initialize()
data = await collector.collect(params.get("sources", []))
self.pipeline_state["data"] = data
elif stage_name == "analyze":
analyst = self.agents["codex-analyst"]
result = await analyst.analyze(
self.pipeline_state["data"],
params.get("analysis_types", [])
)
self.pipeline_state["analysis"] = result
elif stage_name == "report":
reporter = self.agents["claude-reporter"]
report_paths = await reporter.generate_report(
self.pipeline_state["analysis"],
params.get("report_config", {})
)
self.pipeline_state["reports"] = report_paths
elif stage_name == "distribute":
distributor = self.agents["distributor"]
await distributor.distribute(
self.pipeline_state["reports"],
params.get("recipients", [])
)
stage_end = datetime.now()
self.pipeline_state["stages"][stage_name] = {
"duration": (stage_end - stage_start).total_seconds(),
"status": "completed",
}5.2 运行效果示例
🚀 启动管道: monthly-report (2024-01-01 10:00:00)
📥 Stage 1: 数据采集
├── API: sales_metrics ... ✅ 12,458 条记录
├── DB: user_analytics ... ✅ 8,923 条记录
└── Web: competitor_data ... ✅ 156 条记录
⏱️ 耗时: 45s
🧹 Stage 2: 数据清洗
├── 去除重复: 234 条
├── 填充缺失值: 89 条
├── 异常值标记: 45 条
└── 格式标准化: 完成
⏱️ 耗时: 12s
📊 Stage 3: 数据分析
├── 趋势分析: 完成 (7个指标)
├── 异常检测: 45 个异常值
├── 相关性分析: 3 对强相关
└── 洞察生成: 8 条关键发现
⏱️ 耗时: 28s
📝 Stage 4: 报告生成
├── Markdown: ✅
├── HTML: ✅
└── PDF: ✅
⏱️ 耗时: 35s
📤 Stage 5: 报告分发
├── Email: team@example.com ... ✅
└── WeChat: 技术团队群 ... ✅
⏱️ 耗时: 8s
✅ 管道完成! 总耗时: 128s六、监控与运维
6.1 管道健康监控
# config/monitoring.yaml
monitoring:
metrics:
- name: "pipeline_duration"
type: "histogram"
labels: ["pipeline_name"]
- name: "pipeline_success_rate"
type: "counter"
labels: ["pipeline_name", "status"]
- name: "data_freshness"
type: "gauge"
labels: ["source_name"]
alerts:
- name: "pipeline-failure"
condition: "pipeline_success_rate < 0.95"
window: "1h"
notify: ["wechat:admin", "pagerduty"]
- name: "stale-data"
condition: "data_freshness > 3600"
notify: ["slack:#data-alerts"]6.2 故障恢复
# src/pipeline/recovery.py
class PipelineRecovery:
"""管道故障恢复"""
async def recover(self, failed_stage: str,
state: Dict[str, Any]) -> Dict:
"""从失败的阶段恢复"""
recovery_strategy = self._get_recovery_strategy(failed_stage)
if recovery_strategy == "retry":
return await self._retry_stage(failed_stage, state)
elif recovery_strategy == "skip":
return await self._skip_stage(failed_stage, state)
elif recovery_strategy == "fallback":
return await self._use_fallback_data(failed_stage, state)
def _get_recovery_strategy(self, stage: str) -> str:
"""根据阶段决定恢复策略"""
strategies = {
"collect": "retry", # 采集失败 → 重试
"clean": "retry", # 清洗失败 → 重试
"analyze": "fallback", # 分析失败 → 用上次的
"report": "skip", # 报告失败 → 跳过
}
return strategies.get(stage, "retry")总结
通过 Hermes Cron 定时触发 + Codex 数据采集分析 + Claude Code 报告生成,我们构建了一套完整的多 Agent 数据管道系统:
- 全自动运行:无需人工干预,定时自动完成采集 → 分析 → 报告全流程
- 模块化设计:每个阶段独立运行,可随时替换或升级单个 Agent
- 多格式输出:支持 Markdown、HTML、PDF 等多种报告格式
- 多渠道分发:邮件、微信、Slack 一键推送
- 容错恢复:失败自动重试,支持回退到历史数据
- 可观测性:完整的监控指标和告警机制
这套系统适用于:
- 业务数据的日报/周报/月报自动生成
- 运营指标的实时监控与分析
- 竞品数据的定期采集与对比
- 系统日志的自动化分析告警
下篇预告
下一篇 83-混合Provider策略,我们将深入探讨如何在多 Agent 系统中实现成本与质量的最优平衡:简单任务交给本地模型处理、复杂任务路由到 Claude Sonnet、关键结果通过 GPT-4o 审查验证。用最少的成本获得最好的结果,敬请期待!