第 24 章 · 实战二:数据分析流水线
本章目标:用 Flow 事件驱动编排一条"CSV 数据集 → 清洗建议 → 统计分析 → 报告"的流水线,掌握 Pydantic 结构化状态、
@router规模分支、pandas 封装成工具三个生产级技巧。
24.1 为什么这条流水线该用 Flow 而不是 Crew
与第 23 章不同,数据分析有确定性代码环节(读 CSV、算统计量)和分支逻辑(小数据直接分析,大数据先采样)。Crew 擅长角色协作,Flow 擅长精确控制执行路径——所以正确形态是:
@start 读数据(代码) → @router 判断规模(代码)
├─ 小文件 ────────────────→ 分析 crew
└─ 大文件 → 采样 crew ──→ 分析 crew
↓
@listen 汇总生成报告(crew)原则回顾(第 15 章):能用代码确定的绝不用 LLM 决定;LLM 只负责清洗建议、结论解读这类真正的语言任务。
24.2 用 Pydantic 定义 Flow 状态
结构化状态让每个步骤的输入输出都有类型约束,比裸 dict 可靠得多:
# state.py —— Flow 的共享状态定义
from pathlib import Path
from pydantic import BaseModel, Field
class DatasetInfo(BaseModel):
"""步骤 1 读出的数据集元信息"""
path: str
rows: int
cols: int
columns: list[str] = Field(default_factory=list)
missing_ratio: float = 0.0 # 整表缺失率,决定清洗力度
size_class: str = "small" # 由 router 写回:small / large
class CleaningAdvice(BaseModel):
"""采样 crew 产出的清洗建议(结构化输出)"""
issues: list[str] = Field(description="发现的数据质量问题")
recommendations: list[str] = Field(description="逐条清洗建议")
class AnalysisState(BaseModel):
"""整个 Flow 的根状态"""
dataset: DatasetInfo | None = None
cleaning: CleaningAdvice | None = None
stats_summary: str = ""
final_report: str = ""
class Config:
# 允许 pandas 等对象临时进出状态而不被校验卡死
arbitrary_types_allowed = True
def probe_csv(path: str) -> DatasetInfo:
"""纯代码步骤:探测 CSV 元信息,不消耗任何 token"""
import pandas as pd
df = pd.read_csv(path)
return DatasetInfo(
path=str(Path(path).resolve()),
rows=len(df),
cols=len(df.columns),
columns=list(df.columns),
missing_ratio=float(df.isna().mean().mean()),
)probe_csv 是普通 Python 函数——读文件、数行列、算缺失率,零 LLM 参与。这就是 Flow 的价值:把"能算的"从"要问模型的"里剥离出去。
24.3 把 pandas 封装成工具
分析 crew 需要真正"摸到"数据。把统计计算封装成 BaseTool,模型只负责决定调用哪个统计、解读结果:
# tools.py —— 数据分析工具箱(第 8 章模式的实战应用)
import pandas as pd
from crewai.tools import BaseTool
from pydantic import BaseModel, Field
class DescribeInput(BaseModel):
path: str = Field(description="CSV 文件绝对路径")
columns: list[str] = Field(default_factory=list, description="只统计这些列,空则全表")
class DescribeDatasetTool(BaseTool):
name: str = "describe_dataset"
description: str = (
"对 CSV 数值列计算 count/mean/std/min/quartiles/max。"
"用于了解字段分布与异常值。"
)
args_schema: type[BaseModel] = DescribeInput
def _run(self, path: str, columns: list[str] | None = None) -> str:
df = pd.read_csv(path)
if columns:
df = df[columns]
desc = df.describe(include="number").round(3)
return desc.to_string()
class CorrelationInput(BaseModel):
path: str = Field(description="CSV 文件绝对路径")
top_n: int = Field(default=5, description="返回相关性最强的前 N 对数值列")
class CorrelationTool(BaseTool):
name: str = "correlation_matrix"
description: str = "计算数值列皮尔逊相关矩阵,返回相关性最强的前 N 对列及系数"
args_schema: type[BaseModel] = CorrelationInput
def _run(self, path: str, top_n: int = 5) -> str:
corr = pd.read_csv(path).corr(numeric_only=True)
pairs = (
corr.where(lambda x: abs(x) < 0.999) # 去掉自相关 1.0
.stack()
.sort_values(key=abs, ascending=False) # 按相关强度排序
.drop_duplicates()
)
head = pairs.head(top_n * 2).iloc[::2] # 去掉 ±对称重复
return "\n".join(f"{a} ~ {b}: {v:.3f}" for (a, b), v in head.items())封装要点(对照第 8 章):args_schema 用 Pydantic 强约束参数,模型填错会立刻报校验错误;_run 返回紧凑字符串而不是整张 DataFrame——控制进入上下文的 token 量是数据类工具的生命线。
24.4 主流程:@start / @router / @listen 编排
# flow.py —— 数据分析流水线主流程
import os
from crewai import Agent, Crew, Task, LLM
from crewai.flow.flow import Flow, listen, router, start
from pydantic import BaseModel
from state import AnalysisState, CleaningAdvice, probe_csv
from tools import CorrelationTool, DescribeDatasetTool
llm = LLM(
model="openai/deepseek-chat",
base_url="https://api.deepseek.com/v1",
api_key=os.getenv("DEEPSEEK_API_KEY"),
)
ROW_THRESHOLD = 50_000 # 超过 5 万行视为大数据集
def make_cleaning_crew() -> Crew:
advisor = Agent(
role="数据质量顾问",
goal="审阅数据集概况并给出可执行的清洗建议",
backstory="资深数据工程师,见惯脏数据,建议永远具体到列名和操作",
llm=llm,
)
task = Task(
description=(
"数据集 {path} 共 {rows} 行 {cols} 列,整体缺失率 "
"{missing_ratio:.1%}。请评估质量风险并给出清洗建议。"
),
expected_output="问题清单与逐条建议(中文)",
agent=advisor,
output_pydantic=CleaningAdvice, # 结构化输出,直接进状态
)
return Crew(agents=[advisor], tasks=[task])
def make_analysis_crew() -> Crew:
analyst = Agent(
role="数据分析师",
goal="用统计工具完成探索性分析并解读发现",
backstory="十年经验的数据科学家,坚持每个结论都必须由数字支撑",
llm=llm,
tools=[DescribeDatasetTool(), CorrelationTool()],
max_iter=6,
)
task = Task(
description=(
"对 {path} 做探索性分析:先用 describe_dataset 了解分布,"
"再用 correlation_matrix 找强相关对,最后给出 3 条业务解读。"
),
expected_output="关键统计量摘要 + 相关性发现 + 3 条带数字支撑的解读",
agent=analyst,
)
return Crew(agents=[analyst], tasks=[task])
def make_report_crew() -> Crew:
writer = Agent(
role="报告撰写人",
goal="把清洗与分析结果整合成管理层能看懂的报告",
backstory="擅长把技术细节翻译成商业语言的咨询顾问",
llm=llm,
)
task = Task(
description="整合以下素材写一份 500 字内分析报告:\n清洗结论:\n{cleaning}\n\n分析发现:\n{stats}",
expected_output="Markdown 报告:背景、数据质量、核心发现、行动建议",
agent=writer,
)
return Crew(agents=[writer], tasks=[task])
class DataPipeline(Flow[AnalysisState]):
@start()
def load_dataset(self):
"""步骤 1:纯代码探测数据集,写入状态"""
self.state.dataset = probe_csv("data/sales.csv")
print(f"[load] {self.state.dataset.rows} 行 x {self.state.dataset.cols} 列")
@router(load_dataset)
def route_by_size(self):
"""步骤 2:按行数分支——决策权在代码不在 LLM"""
if self.state.dataset.rows > ROW_THRESHOLD:
return "large"
return "small"
@listen("large")
def sample_and_advise(self):
"""大文件分支:先采样再给清洗建议"""
sample_path = "data/_sample.csv"
import pandas as pd
pd.read_csv(self.state.dataset.path, nrows=10_000).to_csv(sample_path, index=False)
advice = make_cleaning_crew().kickoff(
inputs=self.state.dataset.model_dump() | {"path": sample_path}
)
self.state.cleaning = advice.pydantic
@listen("small")
def advise_directly(self):
"""小文件分支:直接给清洗建议"""
advice = make_cleaning_crew().kickoff(inputs=self.state.dataset.model_dump())
self.state.cleaning = advice.pydantic
@listen(sample_and_advise) # or_ 合并两条分支:任一完成即触发(见下方说明)
def analyze(self):
result = make_analysis_crew().kickoff(inputs={"path": self.state.dataset.path})
self.state.stats_summary = result.raw
@listen(analyze)
def build_report(self):
report = make_report_crew().kickoff(inputs={
"cleaning": self.state.cleaning.model_dump_json(indent=2),
"stats": self.state.stats_summary,
})
self.state.final_report = report.raw
print("\n===== 最终报告 =====\n", self.state.final_report)
if __name__ == "__main__":
DataPipeline().kickoff()关于双分支汇合
示例中 analyze 只监听了 sample_and_advise 一条边以保持代码简单;真实实现两条分支汇合时应使用官方的 or_(sample_and_advise, advise_directly) 作为监听目标——or_ 表示任一来源完成即触发,对应的还有 and_(全部完成才触发),二者均来自 crewai.flow.flow。
24.5 运行结果与扩展方向
mkdir -p data && python -c "
import pandas as pd, numpy as np
df = pd.DataFrame({'revenue': np.random.gamma(5, 200, 800),
'visits': np.random.poisson(30, 800)})
df.to_csv('data/sales.csv', index=False)"
python -m flow # 或 crewai run flow典型日志顺序印证了编排逻辑:[load] 800 行 x 2 列 → 走 small 分支 → 清洗建议结构化落库 → 分析师交替调用两个统计工具 → 报告产出。可继续演进的方向:
- checkpoint 持久化:给 Flow 加
@persist或第 22 章 checkpoint,长分析中断可续跑:
from crewai.flow.persistence import persist
class DurablePipeline(DataPipeline):
@persist # 每个方法执行后自动把 state 快照落盘,重跑时恢复到最新快照继续
@listen(analyze)
def build_report(self):
report = make_report_crew().kickoff(inputs={
"cleaning": self.state.cleaning.model_dump_json(indent=2),
"stats": self.state.stats_summary,
})
self.state.final_report = report.raw- 更多分支维度:在
route_by_size里同时按missing_ratio > 40%分流到"重清洗"链路; - 接入 API 化:Flow 同样适用第 22 章的 kickoff+轮询模式,
FlowState天然适合序列化进 job 表。
本章小结
- 选择 Flow 的判据:流程含确定性代码环节或分支路由时,Flow + Crew 组合优于纯 Crew;
- Pydantic 根状态 +
output_pydantic让步骤间传递的是经过校验的结构化对象; - pandas 能力封装成
BaseTool(args_schema校验、返回紧凑字符串),LLM 负责"选哪个统计、如何解读"; @start/@router/@listen三件套完成事件驱动编排,规模阈值等决策留在 Python 代码里;- 双分支汇合用
or_/and_控制触发语义。
🧪 随堂测验
点击你认为正确的选项。答错时会展示正确答案与原因解析。
1. 本项目中"按行数判断走哪条分支"为什么放在 @router 的普通 Python 代码里,而不是让 Agent 决定?
2. Task 设置 output_pydantic=CleaningAdvice 后,kickoff 结果中清洗建议存放在哪里?
3. 把 pandas describe 结果作为工具返回值时,为什么要转成紧凑字符串并限制行数?
4. Flow 中两条互斥分支(large/small)之后都要执行统计分析,正确的监听写法是?
🛠️ 动手实践
- 给
probe_csv增加检测:找出缺失率超过 60% 的列名列表存入DatasetInfo.sparse_columns,并在清洗建议任务的 description 中显式提醒顾问关注这些列。 - 新增一个
OutlierTool工具:基于 IQR(四分位距)规则返回每列的离群值行号与数量上限摘要,让分析 crew 在报告中引用离群值证据。 - 在
route_by_size中增加第三个出口"dirty"(当missing_ratio > 0.4时返回),为其编写一个"深度清洗 crew"分支,并用or_把三条分支汇合到analyze。
最后一个实战:把 Crew 的记忆与定时调度用起来,搭建竞品监控情报系统:第 25 章 · 实战三。