跳转到主要内容

文档索引

获取完整文档索引: https://docs.crewai.com.cn/llms.txt

在深入了解之前,请使用此文件来浏览所有可用页面。

概述

CrewAI Flows 是一项强大的功能,旨在简化 AI 工作流的创建和管理。Flows 允许开发者高效地组合和协调编码任务与智能体团队 (Crews),为构建复杂的 AI 自动化提供了一个健壮的框架。 Flows 使您能够创建结构化的、事件驱动的工作流。它们提供了一种无缝的方式来连接多个任务、管理状态并控制 AI 应用程序中的执行流。通过 Flows,您可以轻松设计和实现充分利用 CrewAI 能力的多步骤流程。
  1. 简化工作流创建:轻松串联多个智能体团队和任务,构建复杂的工作流。
  2. 状态管理:Flows 让您能够极易地在工作流的各个任务之间管理和共享状态。
  3. 事件驱动架构:基于事件驱动模型构建,支持动态且响应迅速的工作流。
  4. 灵活的控制流:在工作流中实现条件逻辑、循环和分支。

开始入门

让我们创建一个简单的流程:在第一个任务中使用 OpenAI 生成一个随机城市,然后在另一个任务中使用该城市生成一个趣闻。
代码

from crewai.flow.flow import Flow, listen, start
from dotenv import load_dotenv
from litellm import completion

load_dotenv()

class ExampleFlow(Flow):
    model = "gpt-4o-mini"

    @start()
    def generate_city(self):
        print("Starting flow")
        # Each flow state automatically gets a unique ID
        print(f"Flow State ID: {self.state['id']}")

        response = completion(
            model=self.model,
            messages=[
                {
                    "role": "user",
                    "content": "Return the name of a random city in the world.",
                },
            ],
        )

        random_city = response["choices"][0]["message"]["content"]
        # Store the city in our state
        self.state["city"] = random_city
        print(f"Random City: {random_city}")

        return random_city

    @listen(generate_city)
    def generate_fun_fact(self, random_city):
        response = completion(
            model=self.model,
            messages=[
                {
                    "role": "user",
                    "content": f"Tell me a fun fact about {random_city}",
                },
            ],
        )

        fun_fact = response["choices"][0]["message"]["content"]
        # Store the fun fact in our state
        self.state["fun_fact"] = fun_fact
        return fun_fact



flow = ExampleFlow()
flow.plot()
result = flow.kickoff()

print(f"Generated fun fact: {result}")
流程可视化图 在上述示例中,我们创建了一个简单的流程,利用 OpenAI 生成一个随机城市,并针对该城市生成趣闻。该流程包含两个任务:generate_citygenerate_fun_factgenerate_city 是流程的起点,而 generate_fun_fact 监听 generate_city 的输出。 每个 Flow 实例的状态中都会自动接收一个唯一标识符 (UUID),这有助于跟踪和管理执行。状态还可以存储额外数据(如生成的城市和趣闻),这些数据在流程执行期间将保持持久化。 当你运行此流程时,它将:
  1. 为流程状态生成一个唯一 ID
  2. 生成一个随机城市并将其存储在状态中
  3. 生成关于该城市的趣闻并将其存储在状态中
  4. 将结果打印到控制台
状态的唯一 ID 和存储的数据对于跟踪工作流执行以及在任务间维持上下文非常有用。 注意:请确保已设置好 .env 文件以存储您的 OPENAI_API_KEY。此密钥对于验证 OpenAI API 请求是必需的。

@start()

@start() 装饰器标记流程的入口点。您可以:
  • 声明多个无条件的起点:@start()
  • 基于先前的方法或路由标签设置启动门控:@start("method_or_label")
  • 提供可调用的条件,以控制何时触发启动
当流程开始或恢复时,所有满足 @start() 条件的方法都将执行(通常并行)。

@listen()

@listen() 装饰器用于将方法标记为流程中另一个任务的输出监听器。使用 @listen() 装饰的方法将在指定的任务发出输出时执行。该方法可以访问它所监听的任务的输出作为参数。

用法

@listen() 装饰器有几种使用方式:
  1. 按名称监听方法:您可以将要监听的方法名称作为字符串传递。当该方法完成时,监听器方法将被触发。
    代码
    @listen("generate_city")
    def generate_fun_fact(self, random_city):
        # Implementation
    
  2. 直接监听方法:您可以直接传递方法本身。当该方法完成时,监听器方法将被触发。
    代码
    @listen(generate_city)
    def generate_fun_fact(self, random_city):
        # Implementation
    

流程输出

访问和处理流程输出对于将您的 AI 工作流集成到更大的应用程序或系统中至关重要。CrewAI Flows 提供了直接的机制来获取最终输出、访问中间结果并管理流程的整体状态。

获取最终输出

当你运行流程时,最终输出由最后完成的方法决定。kickoff() 方法将返回该最终方法的输出。 以下是获取最终输出的方法:
from crewai.flow.flow import Flow, listen, start

class OutputExampleFlow(Flow):
    @start()
    def first_method(self):
        return "Output from first_method"

    @listen(first_method)
    def second_method(self, first_output):
        return f"Second method received: {first_output}"


flow = OutputExampleFlow()
flow.plot("my_flow_plot")
final_output = flow.kickoff()

print("---- Final Output ----")
print(final_output)
流程可视化图 在此示例中,second_method 是最后完成的方法,因此其输出即为流程的最终输出。kickoff() 方法返回该输出,随后将其打印到控制台。plot() 方法将生成 HTML 文件,帮助您理解流程结构。

访问和更新状态

除了获取最终输出,您还可以访问并更新流程内的状态。状态可用于在不同方法间存储和共享数据。流程运行后,您可以访问状态以获取执行期间添加或更新的任何信息。 以下是更新和访问状态的示例:
from crewai.flow.flow import Flow, listen, start
from pydantic import BaseModel

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()
flow.plot("my_flow_plot")
final_output = flow.kickoff()
print(f"Final Output: {final_output}")
print("Final State:")
print(flow.state)
流程可视化图 在此示例中,状态被 first_methodsecond_method 同时更新。流程运行后,您可以访问最终状态来查看这些方法所做的更新。 通过确保返回最终方法的输出并提供状态访问能力,CrewAI Flows 使得将 AI 工作流结果集成到更大的应用程序中变得简单,同时还能在整个流程执行期间维护和访问状态。

流程状态管理

有效管理状态对于构建可靠且可维护的 AI 工作流至关重要。CrewAI Flows 为非结构化和结构化状态管理提供了健壮的机制,使开发者能够选择最适合其应用程序需求的方法。

非结构化状态管理

在非结构化状态管理中,所有状态都存储在 Flow 类的 state 属性中。这种方法提供了灵活性,使开发者无需定义严格的模式即可即时添加或修改状态属性。即使在非结构化状态下,CrewAI Flows 也会自动为每个状态实例生成并维护一个唯一标识符 (UUID)。
代码
from crewai.flow.flow import Flow, listen, start

class UnstructuredExampleFlow(Flow):

    @start()
    def first_method(self):
        # The state automatically includes an 'id' field
        print(f"State ID: {self.state['id']}")
        self.state['counter'] = 0
        self.state['message'] = "Hello from structured flow"

    @listen(first_method)
    def second_method(self):
        self.state['counter'] += 1
        self.state['message'] += " - updated"

    @listen(second_method)
    def third_method(self):
        self.state['counter'] += 1
        self.state['message'] += " - updated again"

        print(f"State after third_method: {self.state}")


flow = UnstructuredExampleFlow()
flow.plot("my_flow_plot")
flow.kickoff()
流程可视化图 注意: id 字段会自动生成并贯穿流程执行始终。您无需手动管理或设置它,即使在更新具有新数据的状态时,它也会保持不变。 关键点:
  • 灵活性: 您可以动态地向 self.state 添加属性,而无需预定义约束。
  • 简单性: 非常适合状态结构较少或变化显著的直接工作流。

结构化状态管理

结构化状态管理利用预定义的模式来确保工作流的一致性和类型安全。通过使用 Pydantic 的 BaseModel 等模型,开发者可以定义状态的确切形状,从而在开发环境中实现更好的验证和自动补全。 CrewAI Flows 中的每个状态都会自动接收一个唯一标识符 (UUID),以帮助跟踪和管理状态实例。此 ID 由 Flow 系统自动生成和管理。
代码
from crewai.flow.flow import Flow, listen, start
from pydantic import BaseModel


class ExampleState(BaseModel):
    # Note: 'id' field is automatically added to all states
    counter: int = 0
    message: str = ""


class StructuredExampleFlow(Flow[ExampleState]):

    @start()
    def first_method(self):
        # Access the auto-generated ID if needed
        print(f"State ID: {self.state.id}")
        self.state.message = "Hello from structured flow"

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

    @listen(second_method)
    def third_method(self):
        self.state.counter += 1
        self.state.message += " - updated again"

        print(f"State after third_method: {self.state}")


flow = StructuredExampleFlow()
flow.kickoff()
流程可视化图 关键点:
  • 定义模式: ExampleState 清晰地勾勒出状态结构,增强了代码的可读性和可维护性。
  • 类型安全: 利用 Pydantic 可确保状态属性符合指定的类型,从而减少运行时错误。
  • 自动补全: IDE 可以基于定义的状态模型提供更好的自动补全和错误检查。

在非结构化和结构化状态管理之间进行选择

  • 在以下情况下使用非结构化状态管理:
    • 工作流的状态简单或高度动态。
    • 比起严格的状态定义,更优先考虑灵活性。
    • 需要快速原型设计,且不想增加定义模式的开销。
  • 在以下情况下使用结构化状态管理:
    • 工作流需要定义明确且一致的状态结构。
    • 类型安全和验证对应用程序的可靠性非常重要。
    • 希望利用 IDE 的自动补全和类型检查功能以获得更好的开发者体验。
通过提供非结构化和结构化状态管理选项,CrewAI Flows 使开发者能够构建既灵活又健壮的 AI 工作流,满足广泛的应用程序需求。

流程持久化

@persist 装饰器可以在 CrewAI Flows 中启用自动状态持久化,允许您在重启或执行不同工作流后维护流程状态。此装饰器既可应用于类级别,也可应用于方法级别,从而在管理状态持久化方面提供了灵活性。

类级别持久化

应用于类级别时,@persist 装饰器会自动持久化所有流程方法的状态。
@persist  # Using SQLiteFlowPersistence by default
class MyFlow(Flow[MyState]):
    @start()
    def initialize_flow(self):
        # This method will automatically have its state persisted
        self.state.counter = 1
        print("Initialized flow. State ID:", self.state.id)

    @listen(initialize_flow)
    def next_step(self):
        # The state (including self.state.id) is automatically reloaded
        self.state.counter += 1
        print("Flow state is persisted. Counter:", self.state.counter)

方法级别持久化

为了实现更精细的控制,您可以将 @persist 应用于特定方法。
class AnotherFlow(Flow[dict]):
    @persist  # Persists only this method's state
    @start()
    def begin(self):
        if "runs" not in self.state:
            self.state["runs"] = 0
        self.state["runs"] += 1
        print("Method-level persisted runs:", self.state["runs"])

分支持久化状态

@persistkickoff / kickoff_async 上支持两种不同的水合 (hydration) 模式:
  • kickoff(inputs={"id": <uuid>})恢复:为提供的 UUID 加载最新的快照,并继续在相同的 flow_uuid 下写入。历史记录会延长。
  • kickoff(restore_from_state_id=<uuid>)分支 (fork):为提供的 UUID 加载最新的快照,从中初始化新运行的状态,并分配一个新的 state.id(自动生成,或在 inputs["id"] 已固定时使用该值)。新运行的 @persist 写入将落在新的 state.id 下;源流程的历史记录得到保留。
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
class CounterFlow(Flow[CounterState]):
    @start()
    def step(self):
        self.state.counter += 1
        print(f"[id={self.state.id}] counter={self.state.counter}")

# Run 1: fresh state, counter 0 -> 1, persisted under flow_1.state.id
flow_1 = CounterFlow()
flow_1.kickoff()

# Fork: hydrate from flow_1's latest snapshot, but use a NEW state.id
flow_2 = CounterFlow()
flow_2.kickoff(restore_from_state_id=flow_1.state.id)
# flow_2.state.counter starts at 1 (hydrated), then step() bumps it to 2.
# flow_2.state.id != flow_1.state.id; flow_1's history is unchanged.
如果提供的 restore_from_state_id 与任何持久化状态不匹配,kickoff 将静默回退 — 这与现有的 inputs["id"] 恢复未找到行为相同。将 restore_from_state_idfrom_checkpoint 组合使用会引发 ValueError;请选择一种水合源。在分支时固定 inputs["id"] 会与其他流程共享持久化密钥 — 通常您只需 restore_from_state_id

工作原理

  1. 唯一状态标识
    • 每个流程状态都会自动接收一个唯一的 UUID
    • ID 在状态更新和方法调用中得到保留
    • 同时支持结构化 (Pydantic BaseModel) 和非结构化 (字典) 状态
  2. 默认 SQLite 后端
    • SQLiteFlowPersistence 是默认存储后端
    • 状态会自动保存到本地 SQLite 数据库
    • 健壮的错误处理确保数据库操作失败时能提供清晰的消息
  3. 错误处理
    • 针对数据库操作的全面错误消息
    • 在保存和加载期间自动进行状态验证
    • 当持久化操作遇到问题时提供清晰反馈

重要注意事项

  • 状态类型:支持结构化 (Pydantic BaseModel) 和非结构化 (字典) 状态
  • 自动 ID:如果不存在,会自动添加 id 字段
  • 状态恢复:失败或重启的流程可以自动重新加载其先前的状态
  • 自定义实现:您可以为特殊存储需求提供自己的 FlowPersistence 实现

技术优势

  1. 通过底层访问实现精确控制
    • 直接访问持久化操作,适用于高级用例
    • 通过方法级持久化装饰器实现精细控制
    • 内置状态检查和调试功能
    • 全面了解状态变化和持久化操作
  2. 增强的可靠性
    • 系统故障或重启后的自动状态恢复
    • 基于事务的状态更新以保证数据完整性
    • 带有清晰错误消息的全面错误处理
    • 在状态保存和加载操作期间进行健壮的验证
  3. 可扩展架构
    • 通过 FlowPersistence 接口可定制持久化后端
    • 支持 SQLite 之外的专门存储解决方案
    • 与结构化 (Pydantic) 和非结构化 (dict) 状态兼容
    • 与现有的 CrewAI 流程模式无缝集成
持久化系统的架构强调技术精度和定制化选项,使开发者能够完全掌控状态管理,同时受益于内置的可靠性功能。

流程控制

条件逻辑:or

Flows 中的 or_ 函数允许您监听多个方法,并在任何指定的方法发出输出时触发监听器方法。
from crewai.flow.flow import Flow, listen, or_, start

class OrExampleFlow(Flow):

    @start()
    def start_method(self):
        return "Hello from the start method"

    @listen(start_method)
    def second_method(self):
        return "Hello from the second method"

    @listen(or_(start_method, second_method))
    def logger(self, result):
        print(f"Logger: {result}")



flow = OrExampleFlow()
flow.plot("my_flow_plot")
flow.kickoff()
流程可视化图 运行此流程时,logger 方法将由 start_methodsecond_method 的输出触发。or_ 函数用于监听多个方法,并在任何一个指定的方法发出输出时触发监听器方法。

条件逻辑:and

Flows 中的 and_ 函数允许您监听多个方法,并仅在所有指定的方法都发出输出时才触发监听器方法。
from crewai.flow.flow import Flow, and_, listen, start

class AndExampleFlow(Flow):

    @start()
    def start_method(self):
        self.state["greeting"] = "Hello from the start method"

    @listen(start_method)
    def second_method(self):
        self.state["joke"] = "What do computers eat? Microchips."

    @listen(and_(start_method, second_method))
    def logger(self):
        print("---- Logger ----")
        print(self.state)

flow = AndExampleFlow()
flow.plot()
flow.kickoff()
流程可视化图 运行此流程时,logger 方法仅在 start_methodsecond_method 都发出输出时才触发。and_ 函数用于监听多个方法,并仅在所有指定的方法都发出输出时触发监听器方法。

路由器 (Router)

Flows 中的 @router() 装饰器允许您根据方法的输出定义条件路由逻辑。您可以根据方法的输出指定不同的路由,从而动态控制执行流。
import random
from crewai.flow.flow import Flow, listen, router, start
from pydantic import BaseModel

class ExampleState(BaseModel):
    success_flag: bool = False

class RouterFlow(Flow[ExampleState]):

    @start()
    def start_method(self):
        print("Starting the structured flow")
        random_boolean = random.choice([True, False])
        self.state.success_flag = random_boolean

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

    @listen("success")
    def third_method(self):
        print("Third method running")

    @listen("failed")
    def fourth_method(self):
        print("Fourth method running")


flow = RouterFlow()
flow.plot("my_flow_plot")
flow.kickoff()
流程可视化图 在上述示例中,start_method 生成一个随机布尔值并将其设置在状态中。second_method 使用 @router() 装饰器根据布尔值定义条件路由逻辑。如果布尔值为 True,该方法返回 "success";如果为 False,则返回 "failed"third_methodfourth_method 监听 second_method 的输出,并根据返回的值执行。 当你运行此流程时,输出将根据 start_method 生成的随机布尔值而变化。

人在回路 (人工反馈)

@human_feedback 装饰器需要 CrewAI 1.8.0 或更高版本
@human_feedback 装饰器通过暂停流程执行以收集人工反馈,从而实现人在回路的工作流。这对于需要人工判断的批准门控、质量审查和决策点非常有用。
代码
from crewai.flow.flow import Flow, start, listen
from crewai.flow.human_feedback import human_feedback, HumanFeedbackResult

class ReviewFlow(Flow):
    @start()
    @human_feedback(
        message="Do you approve this content?",
        emit=["approved", "rejected", "needs_revision"],
        llm="gpt-4o-mini",
        default_outcome="needs_revision",
    )
    def generate_content(self):
        return "Content to be reviewed..."

    @listen("approved")
    def on_approval(self, result: HumanFeedbackResult):
        print(f"Approved! Feedback: {result.feedback}")

    @listen("rejected")
    def on_rejection(self, result: HumanFeedbackResult):
        print(f"Rejected. Reason: {result.feedback}")
当指定 emit 时,人类的自由形式反馈将由 LLM 进行解读并归纳为指定的结果之一,随后触发相应的 @listen 装饰器。 您也可以在不进行路由的情况下使用 @human_feedback 来简单收集反馈:
代码
@start()
@human_feedback(message="Any comments on this output?")
def my_method(self):
    return "Output for review"

@listen(my_method)
def next_step(self, result: HumanFeedbackResult):
    # Access feedback via result.feedback
    # Access original output via result.output
    pass
通过 self.last_human_feedback(最近的)或 self.human_feedback_history(作为列表的所有反馈)访问流程期间收集的所有反馈。 有关流程中人工反馈的完整指南,包括使用自定义提供程序(Slack、Webhooks 等)的 异步/非阻塞反馈,请参阅 流程中的人工反馈

将智能体添加到流程

智能体 (Agents) 可以无缝集成到您的流程中,在您需要更简单、专注的任务执行时,为完整的智能体团队提供了一个轻量级的替代方案。以下是在流程中使用智能体进行市场调研的示例:
import asyncio
from typing import Any, Dict, List

from crewai_tools import SerperDevTool
from pydantic import BaseModel, Field

from crewai.agent import Agent
from crewai.flow.flow import Flow, listen, start


# Define a structured output format
class MarketAnalysis(BaseModel):
    key_trends: List[str] = Field(description="List of identified market trends")
    market_size: str = Field(description="Estimated market size")
    competitors: List[str] = Field(description="Major competitors in the space")


# Define flow state
class MarketResearchState(BaseModel):
    product: str = ""
    analysis: MarketAnalysis | None = None


# Create a flow class
class MarketResearchFlow(Flow[MarketResearchState]):
    @start()
    def initialize_research(self) -> Dict[str, Any]:
        print(f"Starting market research for {self.state.product}")
        return {"product": self.state.product}

    @listen(initialize_research)
    async def analyze_market(self) -> Dict[str, Any]:
        # Create an Agent for market research
        analyst = Agent(
            role="Market Research Analyst",
            goal=f"Analyze the market for {self.state.product}",
            backstory="You are an experienced market analyst with expertise in "
            "identifying market trends and opportunities.",
            tools=[SerperDevTool()],
            verbose=True,
        )

        # Define the research query
        query = f"""
        Research the market for {self.state.product}. Include:
        1. Key market trends
        2. Market size
        3. Major competitors

        Format your response according to the specified structure.
        """

        # Execute the analysis with structured output format
        result = await analyst.kickoff_async(query, response_format=MarketAnalysis)
        if result.pydantic:
            print("result", result.pydantic)
        else:
            print("result", result)

        # Return the analysis to update the state
        return {"analysis": result.pydantic}

    @listen(analyze_market)
    def present_results(self, analysis) -> None:
        print("\nMarket Analysis Results")
        print("=====================")

        if isinstance(analysis, dict):
            # If we got a dict with 'analysis' key, extract the actual analysis object
            market_analysis = analysis.get("analysis")
        else:
            market_analysis = analysis

        if market_analysis and isinstance(market_analysis, MarketAnalysis):
            print("\nKey Market Trends:")
            for trend in market_analysis.key_trends:
                print(f"- {trend}")

            print(f"\nMarket Size: {market_analysis.market_size}")

            print("\nMajor Competitors:")
            for competitor in market_analysis.competitors:
                print(f"- {competitor}")
        else:
            print("No structured analysis data available.")
            print("Raw analysis:", analysis)


# Usage example
async def run_flow():
    flow = MarketResearchFlow()
    flow.plot("MarketResearchFlowPlot")
    result = await flow.kickoff_async(inputs={"product": "AI-powered chatbots"})
    return result


# Run the flow
if __name__ == "__main__":
    asyncio.run(run_flow())
流程可视化图 此示例展示了在流程中使用智能体的几个关键特性:
  1. 结构化输出:使用 Pydantic 模型定义预期输出格式 (MarketAnalysis),确保整个流程的类型安全和结构化数据。
  2. 状态管理:流程状态 (MarketResearchState) 在步骤之间维持上下文,并存储输入和输出。
  3. 工具集成:智能体可以使用工具(如 WebsiteSearchTool)来增强其能力。

将智能体团队添加到流程

在 CrewAI 中创建包含多个智能体团队的流程非常直观。 您可以通过运行以下命令来生成一个新的 CrewAI 项目,其中包含创建多团队流程所需的所有脚手架:
crewai create flow name_of_flow
此命令将生成一个包含必要文件夹结构的新 CrewAI 项目。生成的项目包含一个预构建的名为 poem_crew 的团队,该团队已能直接运行。您可以将此团队用作模板,通过复制、粘贴和编辑它来创建其他团队。

文件夹结构

运行 crewai create flow name_of_flow 命令后,您将看到类似于以下的文件夹结构:
目录/文件描述
name_of_flow/流程的根目录。
├── crews/包含特定团队的目录。
│ └── poem_crew/“poem_crew”团队的目录,包含其配置和脚本。
│ ├── config/“poem_crew”团队的配置文件目录。
│ │ ├── agents.yaml定义“poem_crew”智能体的 YAML 文件。
│ │ └── tasks.yaml定义“poem_crew”任务的 YAML 文件。
│ ├── poem_crew.py“poem_crew”功能的脚本。
├── tools/流程中使用的额外工具目录。
│ └── custom_tool.py自定义工具实现。
├── main.py用于运行流程的主脚本。
├── README.md项目描述和说明。
├── pyproject.toml用于项目依赖项和设置的配置文件。
└── .gitignore指定版本控制中忽略的文件和目录。

构建您的团队

crews 文件夹中,您可以定义多个团队。每个团队都有自己的文件夹,其中包含配置文件和团队定义文件。例如,poem_crew 文件夹包含:
  • config/agents.yaml:定义该团队的智能体。
  • config/tasks.yaml:定义该团队的任务。
  • poem_crew.py:包含团队定义,包括智能体、任务和团队本身。
您可以复制、粘贴并编辑 poem_crew 来创建其他团队。

main.py 中连接团队

main.py 文件是您创建流程并连接各个团队的地方。您可以使用 Flow 类以及 @start@listen 装饰器来定义流程并指定执行顺序。 以下是如何在 main.py 文件中连接 poem_crew 的示例:
代码
#!/usr/bin/env python
from random import randint

from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, start
from .crews.poem_crew.poem_crew import PoemCrew

class PoemState(BaseModel):
    sentence_count: int = 1
    poem: str = ""

class PoemFlow(Flow[PoemState]):

    @start()
    def generate_sentence_count(self):
        print("Generating sentence count")
        self.state.sentence_count = randint(1, 5)

    @listen(generate_sentence_count)
    def generate_poem(self):
        print("Generating poem")
        result = PoemCrew().crew().kickoff(inputs={"sentence_count": self.state.sentence_count})

        print("Poem generated", result.raw)
        self.state.poem = result.raw

    @listen(generate_poem)
    def save_poem(self):
        print("Saving poem")
        with open("poem.txt", "w") as f:
            f.write(self.state.poem)

def kickoff():
    poem_flow = PoemFlow()
    poem_flow.kickoff()


def plot():
    poem_flow = PoemFlow()
    poem_flow.plot("PoemFlowPlot")

if __name__ == "__main__":
    kickoff()
    plot()
在此示例中,PoemFlow 类定义了一个生成句子计数、使用 PoemCrew 生成诗歌,然后将诗歌保存到文件的流程。该流程通过调用 kickoff() 方法启动。PoemFlowPlot 将由 plot() 方法生成。 流程可视化图

运行流程

(可选) 在运行流程之前,您可以通过运行以下命令安装依赖项:
crewai install
安装完所有依赖项后,您需要通过运行以下命令激活虚拟环境:
source .venv/bin/activate
激活虚拟环境后,您可以通过执行以下命令之一来运行流程:
crewai flow kickoff
uv run kickoff
流程将执行,您应该能在控制台中看到输出。

绘制流程图 (Plot Flows)

可视化您的 AI 工作流可以提供关于流程结构和执行路径的宝贵见解。CrewAI 提供了一个强大的可视化工具,允许您生成流程的交互式绘图,使理解和优化您的 AI 工作流变得更加容易。

什么是绘图 (Plots)?

CrewAI 中的绘图是 AI 工作流的图形化表示。它们显示了各种任务、它们的连接以及数据在它们之间的流动。这种可视化有助于理解操作顺序、识别瓶颈,并确保工作流逻辑与您的预期一致。

如何生成绘图

CrewAI 提供了两种方便的方法来生成流程绘图:

选项 1:使用 plot() 方法

如果您直接使用流程实例,可以通过在流程对象上调用 plot() 方法来生成绘图。此方法将创建一个包含流程交互式绘图的 HTML 文件。
代码
# Assuming you have a flow instance
flow.plot("my_flow_plot")
这将在您的当前目录中生成一个名为 my_flow_plot.html 的文件。您可以在网页浏览器中打开此文件以查看交互式绘图。

选项 2:使用命令行

如果您在结构化的 CrewAI 项目中工作,可以使用命令行生成绘图。这对于大型项目特别有用,因为您可以可视化整个流程设置。
crewai flow plot
此命令将生成一个包含流程绘图的 HTML 文件,类似于 plot() 方法。该文件将保存在您的项目目录中,您可以在浏览器中打开它以探索该流程。

理解绘图

生成的绘图将显示表示流程中任务的节点,并带有指示执行流的定向边。绘图是交互式的,允许您进行缩放,并将鼠标悬停在节点上以查看更多详细信息。 通过可视化您的流程,您可以更清楚地理解工作流的结构,从而更容易地调试、优化并向他人传达您的 AI 流程。

结论

绘制流程图是 CrewAI 的一项强大功能,它增强了您设计和管理复杂 AI 工作流的能力。无论您选择使用 plot() 方法还是命令行,生成绘图都将为您提供工作流的视觉表示,有助于开发和演示。

后续步骤

如果您有兴趣探索更多流程示例,我们的示例存储库中有很多推荐内容。这里有四个具体的流程示例,每个示例都展示了独特的用例,帮助您将当前的问题类型与特定示例进行匹配:
  1. 电子邮件自动回复流程:此示例演示了一个无限循环,其中后台作业持续运行以自动回复电子邮件。对于需要无需人工干预而重复执行的任务,这是一个极好的用例。 查看示例
  2. 潜在客户评分流程:此流程展示了如何添加人在回路的反馈,并使用路由器处理不同的条件分支。这是一个极好的示例,说明了如何将动态决策和人工监督整合到您的工作流中。 查看示例
  3. 写书流程:此示例擅长串联多个团队,其中一个团队的输出由另一个团队使用。具体来说,一个团队勾勒整本书的大纲,另一个团队根据大纲生成章节。最终,所有内容连接在一起以产生一本完整的书。此流程非常适合需要不同任务间协调的复杂、多步骤流程。 查看示例
  4. 会议助手流程:此流程演示了如何广播一个事件以触发多个后续操作。例如,会议结束后,流程可以更新 Trello 看板、发送 Slack 消息并保存结果。这是一个处理来自单一事件的多个输出的极好示例,使其成为全面任务管理和通知系统的理想选择。 查看示例
通过探索这些示例,您可以深入了解如何利用 CrewAI Flows 解决各种用例,从自动化重复任务到通过动态决策和人工反馈管理复杂的多步骤流程。 另外,请观看下方关于如何在 CrewAI 中使用流程的 YouTube 视频!

运行流程

有两种运行流程的方法:

使用 Flow API

您可以通过创建流程类的实例并调用 kickoff() 方法以编程方式运行流程。
flow = ExampleFlow()
result = flow.kickoff()

流式传输流程执行

为了实时查看流程执行,您可以启用流式传输,以便在生成输出时接收它。
class StreamingFlow(Flow):
    stream = True  # Enable streaming

    @start()
    def research(self):
        # Your flow implementation
        pass

# Iterate over streaming output
flow = StreamingFlow()
streaming = flow.kickoff()
for chunk in streaming:
    print(chunk.content, end="", flush=True)

# Access final result
result = streaming.result
流式传输流程执行 指南中了解有关流式传输的更多信息。

流程中的记忆 (Memory)

每个流程都会自动访问 CrewAI 的统一 记忆 (Memory) 系统。您可以使用三个内置的便捷方法直接在任何流程方法中存储、回溯和提取记忆。

内置方法

方法描述
self.remember(content, **kwargs)将内容存储在内存中。接受可选的 scopecategoriesmetadataimportance
self.recall(query, **kwargs)检索相关记忆。接受可选的 scopecategorieslimitdepth
self.extract_memories(content)将原始文本分解为离散的、自包含的记忆陈述。
当流程初始化时,默认的 Memory() 实例会自动创建。您也可以传递一个自定义实例。
from crewai.flow.flow import Flow
from crewai import Memory

custom_memory = Memory(
    recency_weight=0.5,
    recency_half_life_days=7,
    embedder={"provider": "ollama", "config": {"model_name": "mxbai-embed-large"}},
)

flow = MyFlow(memory=custom_memory)

示例:研究与分析流程

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


class ResearchAnalysisFlow(Flow):
    @start()
    def gather_data(self):
        # Simulate research findings
        findings = (
            "PostgreSQL handles 10k concurrent connections with connection pooling. "
            "MySQL caps at around 5k. MongoDB scales horizontally but adds complexity."
        )

        # Extract atomic facts and remember each one
        memories = self.extract_memories(findings)
        for mem in memories:
            self.remember(mem, scope="/research/databases")

        return findings

    @listen(gather_data)
    def analyze(self, raw_findings):
        # Recall relevant past research (from this run or previous runs)
        past = self.recall("database performance and scaling", limit=10, depth="shallow")

        context_lines = [f"- {m.record.content}" for m in past]
        context = "\n".join(context_lines) if context_lines else "No prior context."

        return {
            "new_findings": raw_findings,
            "prior_context": context,
            "total_memories": len(past),
        }


flow = ResearchAnalysisFlow()
result = flow.kickoff()
print(result)
由于记忆在运行期间持久存在(由磁盘上的 LanceDB 支持),analyze 步骤也会回溯以前执行中的发现 —— 从而实现随时间学习和积累知识的流程。 有关作用域、切片、复合评分、嵌入器配置等的详细信息,请参阅 记忆文档

使用 CLI

从 0.103.0 版本开始,您可以使用 crewai run 命令运行流程。
crewai run
此命令会自动检测您的项目是否为流程(基于 pyproject.toml 中的 type = "flow" 设置)并相应地运行它。这是从命令行运行流程的推荐方式。 为了向后兼容,您也可以使用:
crewai flow kickoff
然而,crewai run 命令现在是首选方法,因为它适用于智能体团队和流程。