크루 플로우: 이벤트 기반 워크플로우
크루 플로우: 이벤트 기반 워크플로우
CrewAI의 플로우(Flows)는 여러 크루와 태스크를 잇고 상태를 관리하며 실행 흐름을 제어하는 워크플로우 프레임워크예요. 코딩 태스크·크루를 조합해 정교한 AI 자동화를 만들 수 있게 하면서, @start와 @listen 데코레이터로 이벤트 기반 실행을 구성합니다. 여기서는 플로우 생성부터 상태 관리, 조건 분기, 라우터, 크루 연결, 플로팅까지 다룰게요.
출처: 공식문서
본문
개요
CrewAI Flows는 AI 워크플로우의 생성·관리를 간소화하는 강력한 기능입니다. Flows는 코딩 태스크와 Crews를 효율적으로 결합·조정할 수 있게 해 정교한 AI 자동화 구축을 위한 견고한 프레임워크를 제공합니다.
Flows는 구조화되고 이벤트 기반인 워크플로우를 만들 수 있게 해줍니다. 여러 태스크를 연결하고 상태를 관리하며 AI 애플리케이션에서 실행 흐름을 제어하는 원활한 방법을 제공합니다.
- 단순화된 워크플로우 생성: 여러 Crews와 태스크를 쉽게 연결해 복잡한 AI 워크플로우를 구성.
- 상태 관리: 워크플로우의 서로 다른 태스크 간 상태 공유·관리를 쉽게.
- 이벤트 기반 아키텍처: 이벤트 기반 모델 위에 구축되어 역동적·반응형 워크플로우 가능.
- 유연한 제어 흐름: 워크플로우 내 조건 로직, 반복, 분기 구현.
시작하기
한 태스크에서 OpenAI로 랜덤 도시를 만들고, 그 도시로 다른 태스크에서 재미있는 사실을 만드는 간단한 Flow를 만들어 봅시다.
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로 랜덤 도시를 만들고 그 도시에 대한 재미있는 사실을 만드는 간단한 Flow를 만들었습니다. Flow는 generate_city와 generate_fun_fact 두 태스크로 구성됩니다. generate_city가 Flow의 시작점이고, generate_fun_fact는 generate_city의 출력을 리슨합니다.
각 Flow 인스턴스는 상태에 자동으로 고유 식별자(UUID)를 받아 플로우 실행을 추적·관리하는 데 도움이 됩니다. 상태는 생성된 도시·재미있는 사실 같은 추가 데이터도 저장해 플로우 실행 내내 유지됩니다.
상태의 고유 ID와 저장된 데이터는 플로우 실행 추적과 태스크 간 컨텍스트 유지에 유용합니다.
참고: .env 파일에 OPENAI_API_KEY를 저장해 두어야 합니다. 이 키는 OpenAI API 요청 인증에 필요합니다.
@start()
@start() 데코레이터는 Flow의 진입점을 표시합니다.
- 여러 무조건 시작 선언:
@start() - 이전 메서드 또는 라우터 라벨로 시작 게이팅:
@start("method_or_label") - 시작 시점을 제어하는 콜러블 조건 제공
Flow가 시작·재개될 때 모든 충족된 @start() 메서드가 (보통 병렬로) 실행됩니다.
@listen()
@listen() 데코레이터는 Flow의 다른 태스크 출력을 리슨하는 메서드를 표시합니다. 지정된 태스크가 출력을 내면 장식된 메서드가 실행됩니다. 리슨 대상 태스크의 출력을 인자로 받을 수 있습니다.
용법 — @listen() 데코레이터는 여러 방식으로 사용됩니다:
-
메서드 이름으로 리슨: 문자열로 리슨할 메서드 이름 전달. 그 메서드가 완료되면 리스너가 트리거됩니다.
@listen("generate_city") def generate_fun_fact(self, random_city): # Implementation -
메서드 직접 리슨: 메서드 자체를 전달. 완료 시 리스너가 트리거됩니다.
@listen(generate_city) def generate_fun_fact(self, random_city): # Implementation
플로우 출력
최종 출력 얻기
Flow 실행 시 최종 출력은 마지막으로 완료된 메서드가 결정합니다. 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)
---- Final Output ----
Second method received: Output from first_method
이 예시에서 second_method가 마지막으로 완료된 메서드이므로 그 출력이 Flow의 최종 출력이 됩니다. kickoff()가 최종 출력을 반환하고, plot()은 HTML 파일을 생성해 플로우를 이해하는 데 도움을 줍니다.
상태 접근·업데이트
Flow 내에서 상태에 접근하고 업데이트할 수도 있습니다. 상태는 각 메서드 간 데이터를 저장·공유하는 데 쓰입니다. Flow 실행 후 상태에 접근해 실행 중 추가·업데이트된 정보를 얻을 수 있습니다.
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)
Final Output: Hello from first_method - updated by second_method
Final State:
counter=2 message='Hello from first_method - updated by second_method'
플로우 사용 메트릭
Flow 실행이 끝나면 usage_metrics 속성으로 실행 중 모든 LLM 호출의 집계 토큰 사용량을 볼 수 있습니다 — Flow가 오케스트레이션한 각 Crew의 호출, Agent 툴 내부 호출, Flow 메서드의 순수 LLM.call(...) 호출을 포함합니다. 이는 CrewAI Enterprise UI에 표시된 합계의 SDK 측 대응물입니다.
from crewai import LLM
from crewai.flow.flow import Flow, listen, start
class UsageMetricsFlow(Flow):
@start()
def run_first_crew(self):
self.state.first_result = FirstCrew().crew().kickoff()
@listen(run_first_crew)
def call_llm_directly(self):
# Bare LLM call — still counted by flow.usage_metrics
llm = LLM(model="openai/gpt-4o-mini")
self.state.summary = llm.call("Summarize the key takeaways.")
@listen(call_llm_directly)
def run_second_crew(self):
self.state.second_result = SecondCrew().crew().kickoff()
flow = UsageMetricsFlow()
flow.kickoff()
print(flow.usage_metrics)
# UsageMetrics(total_tokens=8579, prompt_tokens=6210, completion_tokens=2369,
# cached_prompt_tokens=0, reasoning_tokens=0,
# cache_creation_tokens=0, successful_requests=5)
UsageMetrics 필드 의미
| 필드 | 의미 |
|---|---|
total_tokens |
청구 총액: prompt_tokens + completion_tokens |
prompt_tokens |
요청에 청구된 전체 입력/프롬프트 토큰 |
completion_tokens |
요청에 청구된 출력/완료 토큰 |
cached_prompt_tokens |
프롬프트 토큰의 캐시 읽기 부분집합(항목만) |
cache_creation_tokens |
프롬프트 토큰의 캐시 쓰기 부분집합(항목만, Anthropic) |
reasoning_tokens |
제공자가 별도로 보고하는 추론/사고 부분집합(항목만) |
successful_requests |
집계된 LLM 호출 수 |
cached_prompt_tokens, cache_creation_tokens, reasoning_tokens 같은 항목 필드는 total_tokens 위에 더해지지 않습니다 — 이들은 이미 prompt_tokens 또는 completion_tokens에 포함된 부분을 설명합니다. Anthropic은 캐시 읽기·쓰기 카운터를 prompt_tokens에 접어 넣어 캐시 워크로드가 total_tokens에 완전히 반영됩니다. OpenAI 스타일 공급자는 캐시 입력을 prompt_tokens 안에 이미 포함하며, CrewAI는 가시성을 위해 캐시된 부분을 별도로 표시합니다.
반환된 UsageMetrics의 각 항목은 단일 flow.kickoff() 호출 내 모든 LLM 호출의 합입니다. 카운터는 다음 kickoff() 호출(또는 kickoff_for_each의 각 반복)에서 리셋되어 연속 실행이 이중 계산되지 않습니다.
플로우 상태 관리
효과적인 상태 관리는 신뢰할 수 있고 유지보수 가능한 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 필드는 자동 생성되어 플로우 실행 내내 유지됩니다. 수동으로 관리·설정할 필요가 없으며, 새 데이터로 상태를 업데이트할 때도 유지됩니다.
구조화 상태 관리
구조화 상태 관리는 사전 정의된 스키마를 활용해 워크플로우 전반의 일관성·타입 안전성을 보장합니다. Pydantic의 BaseModel 같은 모델을 사용해 상태의 정확한 형태를 정의할 수 있어 검증과 개발 환경 자동완성을 개선합니다. 각 상태는 자동으로 고유 UUID를 받습니다.
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()
비구조화 vs 구조화 상태 관리 선택
- 비구조화 상태 관리 사용 시점: 워크플로우 상태가 단순하거나 매우 동적일 때, 엄격한 상태 정의보다 유연성이 우선일 때, 스키마 정의 오버헤드 없이 빠른 프로토타이핑이 필요할 때.
- 구조화 상태 관리 사용 시점: 잘 정의되고 일관된 상태 구조가 필요할 때, 타입 안전성·검증이 애플리케이션 신뢰성에 중요할 때, IDE 자동완성·타입 체크 같은 기능을 활용하고 싶을 때.
플로우 영속화
@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"])
영속 상태 포킹
@persist는 kickoff / kickoff_async에서 두 가지 수화(hydration) 모드를 지원합니다:
kickoff(inputs={"id": <uuid>})— resume: 주어진 UUID의 최신 스냅샷을 로드하고 같은flow_uuid아래 계속 기록. 히스토리가 확장됩니다.kickoff(restore_from_state_id=<uuid>)— fork: 주어진 UUID의 최신 스냅샷을 로드해 새 실행의 상태를 그 스냅샷으로 수화하고 새state.id를 부여(자동 생성 또는inputs["id"]로 고정). 새 실행의@persist쓰기는 새state.id아래 저장되며, 원본 플로우 히스토리는 보존됩니다.
제공된 restore_from_state_id가 어떤 영속 상태와도 일치하지 않으면 kickoff는 조용히 폴백합니다. restore_from_state_id와 from_checkpoint를 함께 쓰면 ValueError가 발생합니다 — 수화 소스는 하나만 선택하세요.
동작 방식
- 고유 상태 식별: 각 플로우 상태는 자동으로 고유 UUID를 받고 상태 업데이트·메서드 호출에 걸쳐 유지됩니다. 구조화(Pydantic BaseModel)·비구조화(딕셔너리) 상태 모두 지원.
- 기본 SQLite 백엔드:
SQLiteFlowPersistence가 기본 저장 백엔드. 상태는 로컬 SQLite DB에 자동 저장. - 오류 처리: DB 연산 실패 시 명확한 오류 메시지, 저장·로드 중 자동 상태 검증.
플로우 제어
조건 로직: or
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: Hello from the start method
Logger: Hello from the second method
조건 로직: and
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 ----
{'greeting': 'Hello from the start method', 'joke': 'What do computers eat? Microchips.'}
라우터
@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()
Starting the structured flow
Third method running
Fourth method running
위 예시에서 start_method가 랜덤 불리언을 생성해 상태에 저장하고, second_method가 @router()로 불리언 값에 따라 조건부 라우팅을 정의합니다. True면 "success", False면 "failed"를 반환하며, third_method와 fourth_method가 그 출력을 리슨해 실행됩니다.
인간 개입(Human in the Loop)
@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(message="...")만 사용할 수도 있습니다. 플로우에서 수집된 모든 피드백은 self.last_human_feedback(가장 최근) 또는 self.human_feedback_history(전체 목록)로 접근할 수 있습니다.
플로우에 에이전트 추가
에이전트는 플로우에 매끄럽게 통합될 수 있어, 더 단순·집중된 태스크 실행이 필요할 때 전체 Crew보다 가벼운 대안이 됩니다. 플로우 내에서 에이전트의 kickoff_async(query, response_format=...)를 호출하고 Pydantic 모델로 구조화 출력을 정의할 수 있습니다.
플로우에 크루 추가
여러 크루가 있는 플로우 생성은 간단합니다. 다음 명령으로 새 CrewAI 프로젝트를 생성하면 됩니다:
crewai create flow name_of_flow
생성된 프로젝트는 poem_crew라는 이미 동작하는 크루를 포함합니다. 폴더 구조:
| 디렉터리/파일 | 설명 |
|---|---|
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 |
버전 관리에서 제외할 파일 지정 |
JSON-first 임베디드 크루는 crew.jsonc와 agents/*.jsonc를 가진 폴더를 사용합니다:
crews/
└── research_crew/
├── agents/
│ └── researcher.jsonc
└── crew.jsonc
그런 다음 Flow 스텝에서 로드합니다:
from pathlib import Path
from crewai.project import load_crew
crew, default_inputs = load_crew(
Path(__file__).parent / "crews" / "research_crew" / "crew.jsonc"
)
result = crew.kickoff(inputs={**default_inputs, "topic": "AI Agents"})
main.py에서 Flow 클래스와 @start/@listen 데코레이터로 크루를 연결합니다. 예를 들어 PoemCrew().crew().kickoff(inputs={...})를 플로우 메서드 안에서 호출해 시를 생성하고 파일에 저장할 수 있습니다.
플로우 실행
의존성 설치: crewai install → 가상환경 활성화: source .venv/bin/activate → 실행: crewai run 또는 uv run kickoff.
플로우 플롯
CrewAI는 플로우의 인터랙티브 플롯을 생성하는 시각화 툴을 제공합니다. flow.plot("my_flow_plot") 메서드로 my_flow_plot.html 파일을 생성하거나, 구조화 프로젝트에서 crewai flow plot 명령으로 생성할 수 있습니다. 플롯은 플로우의 태스크(노드)와 실행 흐름(방향 엣지)을 보여줍니다.
플로우 실행
- Flow API 사용:
flow = ExampleFlow(),result = flow.kickoff(). - 스트리밍 실행:
stream = True로 설정해 출력이 생성되는 대로 받기. - CLI 사용: 버전 0.103.0부터
crewai run으로 플로우를 실행할 수 있습니다(pyproject.toml의type = "flow"설정 자동 감지). 레거시crewai flow kickoff명령은 폐기되었습니다.
플로우에서 메모리
모든 Flow는 CrewAI의 통합 Memory 시스템에 자동 접근합니다. 세 가지 내장 편의 메서드:
| 메서드 | 설명 |
|---|---|
self.remember(content, **kwargs) |
콘텐츠를 메모리에 저장. scope, categories, metadata, importance 선택 인자 |
self.recall(query, **kwargs) |
관련 메모리 검색. scope, categories, limit, depth 선택 인자 |
self.extract_memories(content) |
원시 텍스트를 개별·자족적 메모리 문장으로 분해 |
Flow가 초기화될 때 기본 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)
메모리는 실행 간에 (디스크의 LanceDB로 백업되어) 유지되므로 이전 실행의 발견 내용도 회상할 수 있어, 시간이 지나며 학습·지식을 축적하는 플로우를 만들 수 있습니다.
다음 단계
플로우의 추가 예시는 examples 저장소에서 확인할 수 있습니다: Email Auto Responder Flow(무한 루프 자동 응답), Lead Score Flow(인간-인-더-루프 + 라우터 조건 분기), Write a Book Flow(여러 크루 체이닝), Meeting Assistant Flow(단일 이벤트 → 여러 후속 작업) 등이 대표 사례입니다.