Skip to content

第 24 章 · 实战二:数据分析流水线

本章目标:用 Flow 事件驱动编排一条"CSV 数据集 → 清洗建议 → 统计分析 → 报告"的流水线,掌握 Pydantic 结构化状态、@router 规模分支、pandas 封装成工具三个生产级技巧。

24.1 为什么这条流水线该用 Flow 而不是 Crew

与第 23 章不同,数据分析有确定性代码环节(读 CSV、算统计量)和分支逻辑(小数据直接分析,大数据先采样)。Crew 擅长角色协作,Flow 擅长精确控制执行路径——所以正确形态是:

text
@start 读数据(代码) → @router 判断规模(代码)
   ├─ 小文件 ────────────────→ 分析 crew
   └─ 大文件 → 采样 crew ──→ 分析 crew

                        @listen 汇总生成报告(crew)

原则回顾(第 15 章):能用代码确定的绝不用 LLM 决定;LLM 只负责清洗建议、结论解读这类真正的语言任务。

24.2 用 Pydantic 定义 Flow 状态

结构化状态让每个步骤的输入输出都有类型约束,比裸 dict 可靠得多:

python
# 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,模型只负责决定调用哪个统计、解读结果:

python
# 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 编排

python
# 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 运行结果与扩展方向

bash
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,长分析中断可续跑:
python
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 能力封装成 BaseToolargs_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)之后都要执行统计分析,正确的监听写法是?

🛠️ 动手实践

  1. probe_csv 增加检测:找出缺失率超过 60% 的列名列表存入 DatasetInfo.sparse_columns,并在清洗建议任务的 description 中显式提醒顾问关注这些列。
  2. 新增一个 OutlierTool 工具:基于 IQR(四分位距)规则返回每列的离群值行号与数量上限摘要,让分析 crew 在报告中引用离群值证据。
  3. route_by_size 中增加第三个出口 "dirty"(当 missing_ratio > 0.4 时返回),为其编写一个"深度清洗 crew"分支,并用 or_ 把三条分支汇合到 analyze

最后一个实战:把 Crew 的记忆与定时调度用起来,搭建竞品监控情报系统:第 25 章 · 实战三