Flows 工作流

Flow 是生产应用的骨架:事件驱动、带状态、可分支。官方确认的核心装饰器是 @start()@listen(...)@router();组合条件用 or_ / and_;持久化用 @persistcrewai.flow.persistence)。人机门闸是 @human_feedback(≥ 1.8.0,crewai.flow.human_feedback)。不要发明未出现在文档中的装饰器名。

导入(概念文档):

from crewai.flow.flow import Flow, listen, start, router, or_, and_

Quickstart 也使用 from crewai.flow import Flow, listen, start,二者均可。


状态与入口

每个 Flow 实例的 state 带唯一 id。无结构时用 dict;推荐 Flow[YourModel] + Pydantic。

from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, start

class ExampleState(BaseModel):
    counter: int = 0
    message: str = ""

class StateExampleFlow(Flow[ExampleState]):
    @start()
    def first_method(self):
        self.state.message = "Hello from first_method"
        self.state.counter += 1

    @listen(first_method)
    def second_method(self):
        self.state.message += " - updated by second_method"
        self.state.counter += 1
        return self.state.message

flow = StateExampleFlow()
print(flow.kickoff())
print(flow.state)
  • @start():入口;可以有多个,满足条件时会启动(常并行)。
  • @listen(method)@listen("method_name"):被监听方法完成后运行;可接收其返回值作为参数。
  • kickoff() 返回最后完成的方法的返回值。plot() / plot("name") 生成 HTML 图。

在步骤里跑 Crew

这是官方推荐的组合:Flow 管主题与产物,Crew 做自主研究。

from crewai import Agent, Task, Crew, Process
from crewai.flow.flow import Flow, listen, start

class ResearchFlow(Flow):
    @start()
    def prepare(self):
        self.state["topic"] = "AI Agents"

    @listen(prepare)
    def run_crew(self):
        agent = Agent(
            role="Researcher",
            goal="Brief the topic",
            backstory="Concise, sourced writing.",
        )
        task = Task(
            description="Write a briefing on {topic}.",
            expected_output="Markdown briefing.",
            agent=agent,
            output_file="output/report.md",
        )
        crew = Crew(agents=[agent], tasks=[task], process=Process.sequential)
        result = crew.kickoff(inputs={"topic": self.state["topic"]})
        self.state["report"] = result.raw
        return result.raw

CLI 项目里用 load_crew(Path("crew.jsonc")) 代替手写 Crew(...)。跑完后 flow.usage_metrics 汇总本轮所有 LLM 调用(含 Crew 与 Flow 内直接 LLM.call)。


@routeror_and_

from crewai.flow.flow import Flow, listen, router, start
from pydantic import BaseModel

class FlagState(BaseModel):
    success_flag: bool = False

class RouterFlow(Flow[FlagState]):
    @start()
    def start_method(self):
        self.state.success_flag = True

    @router(start_method)
    def second_method(self):
        return "success" if self.state.success_flag else "failed"

    @listen("success")
    def third_method(self):
        print("ok path")

    @listen("failed")
    def fourth_method(self):
        print("fail path")

@listen(or_(a, b)):任一完成即触发。@listen(and_(a, b)):全部完成才触发。


@persist

from crewai.flow.flow import Flow, start
from crewai.flow.persistence import persist
from pydantic import BaseModel

class CounterState(BaseModel):
    id: str = ""
    counter: int = 0

@persist  # 默认 SQLiteFlowPersistence
class CounterFlow(Flow[CounterState]):
    @start()
    def step(self):
        self.state.counter += 1

kickoff(inputs={"id": ...}) 按同一 UUID 续跑kickoff(restore_from_state_id=...) 分叉出新 state.id。不要和 from_checkpoint 混用。也可把 @persist 只打在某个方法上。


与 LangGraph 怎么选

LangGraph(本站 LangChain 课)用显式 StateGraph、检查点、中断。CrewAI Flow 用装饰器事件图,并原生嵌入 Crew 角色团队。已有 LangGraph 项目可参考官方 Moving from LangGraph to CrewAI;新项目若核心是角色协作,优先 Flow + Crew。


下一步

评论