Skip to content

高级 Flow:分支与路由器

在第 7 章中,我们介绍了流(Flows)作为线性链式事件的方式。然而,实际问题很少是直线式的;它们是决策树。你可能需要先检查某个话题是否敏感,然后再撰写相关内容,或者你可能希望同时生成文本和代码。

在本章中,我们将掌握用于条件逻辑的 @router 装饰器,并探索如何在流(Flow)中并行运行任务。

路由器(Router)允许你的流(Flow)根据当前状态动态决定要走的路径。把它想象成站在十字路口的交通管制员。与 @listen 不同,你可以在一个方法上使用 @router,该方法返回下一个要执行的方法的名称。

想象一个生成博客文章的流。在撰写之前,我们希望检查主题是否与技术相关。如果是,我们编写代码示例;如果不是,我们撰写一篇散文。

from crewai.flow.flow import Flow, start, listen, router
from pydantic import BaseModel
class ContentState(BaseModel):
topic: str = ""
is_tech: bool = False
final_content: str = ""
class ContentFlow(Flow[ContentState]):
@start()
def analyze_topic(self):
# 判断话题是否为技术相关的逻辑
# 在此示例中,我们模拟一个检查
self.state.topic = "Python Decorators"
self.state.is_tech = True
print(f"Analyzing topic: {self.state.topic}")
@router(analyze_topic)
def route_based_on_type(self):
if self.state.is_tech:
return "write_tech_tutorial"
else:
return "write_general_essay"
@listen("write_tech_tutorial")
def write_tech_tutorial(self):
self.state.final_content = "Here is a Python code example..."
print("Writing technical content...")
@listen("write_general_essay")
def write_general_essay(self):
self.state.final_content = "Once upon a time..."
print("Writing essay...")

在上面的代码中,route_based_on_type 本身不执行任务;它严格地引导流的走向。它返回一个与所需下一步方法名称匹配的字符串。

有时效率是关键。如果你需要研究三个不同的网站,逐个进行会很慢。流(Flows)允许你从一个单一的起点触发多个方法。

如果多个方法监听相同的输出(或起点),它们会并行运行。要将它们重新聚合,你可以在 @listen 装饰器内部使用逻辑。

  • @listen(or_=[method_a, method_b]):只要 A 或 B 完成,就立即运行。
  • @listen(and_=[method_a, method_b]):等待 A 和 B 都 完成后才运行。
@start()
def kickoff(self):
print("Starting research...")
@listen(kickoff)
def research_competitor_A(self):
# 模拟工作
return "Data A"
@listen(kickoff)
def research_competitor_B(self):
# 模拟工作
return "Data B"
@listen(and_=[research_competitor_A, research_competitor_B])
def compile_report(self):
# 这只有在两个研究任务都完成后才会运行
print("Compiling final report from A and B...")

通过结合路由器(Routers)和并行执行,你可以构建非线性的、复杂的 AI 流水线,模拟复杂的人类工作流程。