Skip to content

第 12 章 · 输出处理:JSON/Pydantic 与数据传递

本章目标:掌握 output_pydantic / output_json 结构化输出,学会用 CrewOutput.tasks_output 取回每个任务的产物,理解任务间数据的隐式流动与显式传递,并能用 usage_metrics 建立成本账本、搭建输出后处理管道。

12.1 为什么需要结构化输出

crew.kickoff() 默认返回的 raw 是一段自由文本——人看着舒服,程序没法可靠解析。LLM 输出的 Markdown 标题层级、列表符号每次都可能不一样,直接 split() 解析是脆弱的。结构化输出让 crew 的产出直接变成带类型校验的 Python 对象

12.2 output_pydantic 与 output_json

python
import os

from crewai import Agent, Crew, Process, Task, LLM
from pydantic import BaseModel, Field

llm = LLM(
    model="openai/deepseek-chat",
    base_url="https://api.deepseek.com/v1",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
    temperature=0.2,     # 结构化抽取建议低温
)

analyst = Agent(
    role="财报分析师",
    goal="从财报文本中精确提取关键指标",
    backstory="严谨的财务数据专家,只输出事实。",
    llm=llm,
)


# 用 Pydantic 定义期望的输出结构
class FinancialSummary(BaseModel):
    company: str = Field(..., description="公司名称")
    revenue_yi: float = Field(..., description="营业收入(亿元)")
    yoy_growth: float = Field(..., description="营收同比增速(小数,如 0.15)")
    highlights: list[str] = Field(..., description="3 条以内经营亮点")


task = Task(
    description="阅读以下财报摘要并提取关键指标:\n{report_text}",
    expected_output="符合 schema 的结构化数据",
    agent=analyst,
    output_pydantic=FinancialSummary,   # 关键参数
)

crew = Crew(agents=[analyst], tasks=[task], process=Process.sequential)
result = crew.kickoff(inputs={"report_text": "某公司 2025 年营收 120 亿元,同比增长 18%……"})

# result.pydantic 是经过校验的模型实例,可直接当 Python 对象用
summary: FinancialSummary = result.pydantic
print(summary.company)          # 属性访问,无需解析字符串
print(f"增速 {summary.yoy_growth:.1%}")

# 也可以拿 dict / JSON 字符串
print(result.json_dict)         # dict 形式(配置了结构化输出时才有值)

output_json 的用法类似(传入 Pydantic 模型或 dict 结构),区别在于产出以 JSON dict 形式暴露在 json_dict 中;需要类型校验和方法复用时优先选 output_pydantic

只有配置了才有值

官方文档明确:TaskOutput.raw 永远有值,但 pydanticjson_dict 只有当任务配置了对应参数时才非空。下游代码务必做空值防御。

12.3 tasks_output:取回每个任务的产物

CrewOutput 不只有最终结果——tasks_output 列表保留了每个任务的独立产物:

python
result = crew.kickoff()

# 最终产出 = 最后一个任务的输出
print("最终结果:", result.raw)

# 逐个检查中间任务:调试流水线的关键手段
for i, out in enumerate(result.tasks_output, 1):
    print(f"[任务 {i}] agent={out.agent}")
    print(f"  raw 前 60 字: {out.raw[:60]}...")
    if out.pydantic:                       # 该任务若配置了 output_pydantic
        print(f"  结构化字段: {list(out.pydantic.model_fields)}")

# token 用量统计(详见 12.5)
print(crew.usage_metrics)

也可以通过 crew.state["task_outputs"] 之外的常规途径——最常用还是上面遍历 tasks_output 的写法。中间任务同样可以各自配 output_pydantic,实现"每一步都是强类型"的流水线。

12.4 任务间数据传递:隐式流动 vs 显式传递

sequential 流程下,前序任务的输出会自动成为后续任务的上下文(第 5 章讲过),这是隐式传递

python
t1 = Task(description="调研竞品定价策略", expected_output="定价要点列表", agent=researcher)
t2 = Task(
    description="基于调研结果给出我方定价建议",   # "调研结果"自动可见
    expected_output="定价建议表",
    agent=strategist,
)

隐式传递省事,但有两个坑:上下文越长 token 越贵;且 t2 能看到 t1 的全部原文而非提炼版。显式传递context 精确控制信息流:

python
# context 显式声明依赖:只有指定的任务输出会进入当前任务上下文
report_task = Task(
    description="汇总分析与定价结论写成管理层报告",
    expected_output="一页纸决策报告",
    agent=writer,
    context=[analysis_task, pricing_task],   # 只喂这两个,不喂原始调研全文
)

再进一步,可以组合出"结构化接力"模式——上游产 Pydantic 对象,下游任务描述中引用其字段:

python
class PricingAdvice(BaseModel):
    suggested_price: float
    rationale: str

pricing_task = Task(
    description="给出产品定价建议",
    expected_output="结构化定价建议",
    agent=strategist,
    output_pydantic=PricingAdvice,
)

# 在下游 description 中插值上游的结构化结论(kickoff 后由框架渲染 inputs)
launch_task = Task(
    description="根据定价建议 {advice} 起草发布文案",
    expected_output="100 字以内发布文案",
    agent=copywriter,
)
crew.kickoff(inputs={"advice": f"{pricing_task.output.raw}"})

12.5 usage_metrics:给 AI 流水线记账

python
crew.kickoff()
metrics = crew.usage_metrics
print(metrics)
# 包含成功请求次数与 token 统计:
# total_tokens 为计费口径 = prompt_tokens + completion_tokens
# (cached_prompt_tokens 等细分字段已包含在 total 内,不再额外累加)

生产实践中把 usage_metrics 记入监控(配合第 21 章的可观测性方案),按天/按团队聚合,才能在成本失控前发现问题。层级流程(第 9 章)的开销评估尤其依赖这个数字。

12.6 输出后处理管道

推荐的项目分层:crew 负责"生成",普通 Python 代码负责"加工与落库",两者解耦:

python
from datetime import datetime
import json
from pathlib import Path


def post_process(summary: FinancialSummary) -> dict:
    """对结构化产出做业务校验与富化,失败抛异常而不是静默放过。"""
    if summary.yoy_growth < -0.5:
        raise ValueError(f"{summary.company} 增速异常:{summary.yoy_growth:.0%},请人工复核")
    return {
        **summary.model_dump(),
        "captured_at": datetime.now().isoformat(),
        "reviewed": False,
    }


result = crew.kickoff(inputs={"report_text": ...})
record = post_process(result.pydantic)      # 校验+富化
Path("out").mkdir(exist_ok=True)
Path(f"out/{record['company']}.json").write_text(
    json.dumps(record, ensure_ascii=False, indent=2)
)

要点:Pydantic 校验挡住"格式不对",业务校验(如增速阈值)挡住"格式对但数值离谱";两层都过了才允许落库或进入下游系统。

本章小结

  • Task(output_pydantic=Model) 让产出变成带校验的 Python 对象,output_json 则提供 dict 形式;
  • pydantic/json_dict 仅在任务配置了对应参数时才有值,下游要做空值防御;
  • CrewOutput.tasks_output 可逐个取回每个任务的产物,是多步流水线的调试利器;
  • sequential 默认全量上下文流动,context=[...] 可显式收窄信息流,省钱又降噪;
  • usage_metrics 提供计费口径的 token 统计(total = prompt + completion);
  • 后处理应放在 crew 之外:格式校验 + 业务校验双闸门后再落库。

🧪 随堂测验

点击你认为正确的选项。答错时会展示正确答案与原因解析。

1. 任务配置了 output_pydantic=FinancialSummary 后,如何拿到校验后的对象?

2. 某个任务没有配置任何结构化输出参数,它的 TaskOutput 中 json_dict 的值是?

3. 关于任务间数据传递,context=[a_task, b_task] 的作用是?

4. 关于 usage_metrics,正确的说法是?

🛠️ 动手实践

  1. 给一条三任务流水线的每个任务都配上不同的 output_pydantic 模型,跑通后遍历 tasks_output 打印每步的结构化字段。
  2. 把一个 5 任务顺序 crew 改造成用 context 收窄信息流的版本,对比改造前后 usage_metrics 的 token 差异。
  3. 实现 12.6 的后处理管道:为你的场景设计至少两条业务校验规则(如数值阈值、必填字段),并故意构造一次违规输入验证拦截。

至此单条 crew 的"内功"已经学完。下一章进入质量保障:第 13 章 · 规划 Planning、迭代与错误处理