文档索引
获取完整文档索引: https://docs.crewai.com.cn/llms.txt
在深入了解之前,请使用此文件来浏览所有可用页面。
@human_feedback 装饰器要求 CrewAI 版本 1.8.0 或更高版本。在使用此功能前,请确保更新您的安装环境。
@human_feedback 装饰器实现了 CrewAI Flow 内部的“人在回路”(HITL) 工作流。它允许您暂停 Flow 的执行,将输出呈现给人工进行审核,收集其反馈,并根据反馈结果选择性地路由到不同的监听器。 这在以下场景中尤为有价值:
- 质量保证:在 AI 生成的内容用于后续流程前进行审核
- 决策门:在自动化工作流中让人员做出关键决策
- 审批工作流:实施批准/拒绝/修改模式
- 交互式优化:收集反馈以迭代改进输出
快速入门
以下是将人工反馈添加到 Flow 中的最简单方法
from crewai.flow.flow import Flow, start, listen
from crewai.flow.human_feedback import human_feedback
class SimpleReviewFlow(Flow):
@start()
@human_feedback(message="Please review this content:")
def generate_content(self):
return "This is AI-generated content that needs review."
@listen(generate_content)
def process_feedback(self, result):
print(f"Content: {result.output}")
print(f"Human said: {result.feedback}")
flow = SimpleReviewFlow()
flow.kickoff()
当此 Flow 运行时,它将:
- 执行
generate_content 并返回字符串
- 向用户显示输出以及请求消息
- 等待用户输入反馈(或按回车键跳过)
- 将
HumanFeedbackResult 对象传递给 process_feedback
@human_feedback 装饰器
| 参数 | 类型 | 必填 | 描述 |
|---|
message | str | 是 | 与方法输出一起显示给人员的消息 |
emit | Sequence[str] | 否 | 可能的结果列表。反馈会缩减为其中之一,从而触发 @listen 装饰器 |
llm | str | BaseLLM | 指定 emit 时 | 用于解释反馈并映射到结果的 LLM |
default_outcome | str | 否 | 未提供反馈时使用的结果。必须在 emit 中 |
metadata | dict | 否 | 企业集成的额外数据 |
provider | HumanFeedbackProvider | 否 | 用于异步/非阻塞反馈的自定义提供程序。请参阅 异步人工反馈 |
learn | bool | 否 | 启用 HITL 学习:从反馈中提炼经验并预审核未来的输出。默认 False。请参阅 从反馈中学习 |
learn_limit | int | 否 | 预审核时可回顾的最大历史经验数。默认 5 |
基本用法(无路由)
如果不指定 emit,装饰器只会收集反馈并将 HumanFeedbackResult 传递给下一个监听器
@start()
@human_feedback(message="What do you think of this analysis?")
def analyze_data(self):
return "Analysis results: Revenue up 15%, costs down 8%"
@listen(analyze_data)
def handle_feedback(self, result):
# result is a HumanFeedbackResult
print(f"Analysis: {result.output}")
print(f"Feedback: {result.feedback}")
带 emit 的路由
当您指定 emit 时,装饰器将充当路由器。人员的自由文本反馈由 LLM 解释并缩减为指定的其中一个结果
from crewai.flow.flow import Flow, start, listen, or_
from crewai.flow.human_feedback import human_feedback
class ReviewFlow(Flow):
@start()
def generate_content(self):
return "Draft blog post content here..."
@human_feedback(
message="Do you approve this content for publication?",
emit=["approved", "rejected", "needs_revision"],
llm="gpt-4o-mini",
default_outcome="needs_revision",
)
@listen(or_("generate_content", "needs_revision"))
def review_content(self):
return "Draft blog post content here..."
@listen("approved")
def publish(self, result):
print(f"Publishing! User said: {result.feedback}")
@listen("rejected")
def discard(self, result):
print(f"Discarding. Reason: {result.feedback}")
当人员输入类似“需要更多细节”的内容时,LLM 会将其缩减为 "needs_revision",这会通过 or_() 再次触发 review_content —— 从而创建修改循环。循环会持续到结果为 "approved" 或 "rejected" 为止。
LLM 在可用时会使用结构化输出(函数调用)来确保响应是您指定的结果之一。这使得路由变得可靠且可预测。
@start() 方法在 Flow 开始时仅运行一次。如果您需要修改循环,请将起始方法与审核方法分开,并在审核方法上使用 @listen(or_("trigger", "revision_outcome")) 来启用自循环。
HumanFeedbackResult
HumanFeedbackResult 数据类包含有关人工反馈交互的所有信息
from crewai.flow.human_feedback import HumanFeedbackResult
@dataclass
class HumanFeedbackResult:
output: Any # The original method output shown to the human
feedback: str # The raw feedback text from the human
outcome: str | None # The collapsed outcome (if emit was specified)
timestamp: datetime # When the feedback was received
method_name: str # Name of the decorated method
metadata: dict # Any metadata passed to the decorator
在监听器中访问
当监听器被带有 emit 的 @human_feedback 方法触发时,它会接收 HumanFeedbackResult
@listen("approved")
def on_approval(self, result: HumanFeedbackResult):
print(f"Original output: {result.output}")
print(f"User feedback: {result.feedback}")
print(f"Outcome: {result.outcome}") # "approved"
print(f"Received at: {result.timestamp}")
访问反馈历史
Flow 类提供了两个用于访问人工反馈的属性
last_human_feedback
返回最近的 HumanFeedbackResult
@listen(some_method)
def check_feedback(self):
if self.last_human_feedback:
print(f"Last feedback: {self.last_human_feedback.feedback}")
human_feedback_history
Flow 运行期间收集的所有 HumanFeedbackResult 对象列表
@listen(final_step)
def summarize(self):
print(f"Total feedback collected: {len(self.human_feedback_history)}")
for i, fb in enumerate(self.human_feedback_history):
print(f"{i+1}. {fb.method_name}: {fb.outcome or 'no routing'}")
每个 HumanFeedbackResult 都会被追加到 human_feedback_history 中,因此多次反馈步骤不会相互覆盖。使用此列表访问 Flow 期间收集的所有反馈。
完整示例:内容审批工作流
这是一个实现内容审核和审批工作流(包含修改循环)的完整示例
from crewai.flow.flow import Flow, start, listen, or_
from crewai.flow.human_feedback import human_feedback, HumanFeedbackResult
from pydantic import BaseModel
class ContentState(BaseModel):
draft: str = ""
revision_count: int = 0
status: str = "pending"
class ContentApprovalFlow(Flow[ContentState]):
"""A flow that generates content and loops until the human approves."""
@start()
def generate_draft(self):
self.state.draft = "# AI Safety\n\nThis is a draft about AI Safety..."
return self.state.draft
@human_feedback(
message="Please review this draft. Approve, reject, or describe what needs changing:",
emit=["approved", "rejected", "needs_revision"],
llm="gpt-4o-mini",
default_outcome="needs_revision",
)
@listen(or_("generate_draft", "needs_revision"))
def review_draft(self):
self.state.revision_count += 1
return f"{self.state.draft} (v{self.state.revision_count})"
@listen("approved")
def publish_content(self, result: HumanFeedbackResult):
self.state.status = "published"
print(f"Content approved and published! Reviewer said: {result.feedback}")
return "published"
@listen("rejected")
def handle_rejection(self, result: HumanFeedbackResult):
self.state.status = "rejected"
print(f"Content rejected. Reason: {result.feedback}")
return "rejected"
flow = ContentApprovalFlow()
result = flow.kickoff()
print(f"\nFlow completed. Status: {flow.state.status}, Reviews: {flow.state.revision_count}")
关键模式是 @listen(or_("generate_draft", "needs_revision")) —— 审核方法既监听初始触发器,也监听其自身的修改结果,从而创建一个重复循环,直到人员批准或拒绝。
与其他装饰器组合
@human_feedback 装饰器可与 @start()、@listen() 和 or_() 一起使用。两种装饰器顺序都可以工作——框架会在两个方向上传播属性——但推荐的模式是
# One-shot review at the start of a flow (no self-loop)
@start()
@human_feedback(message="Review this:", emit=["approved", "rejected"], llm="gpt-4o-mini")
def my_start_method(self):
return "content"
# Linear review on a listener (no self-loop)
@listen(other_method)
@human_feedback(message="Review this too:", emit=["good", "bad"], llm="gpt-4o-mini")
def my_listener(self, data):
return f"processed: {data}"
# Self-loop: review that can loop back for revisions
@human_feedback(message="Approve or revise?", emit=["approved", "revise"], llm="gpt-4o-mini")
@listen(or_("upstream_method", "revise"))
def review_with_loop(self):
return "content for review"
自循环模式
要创建修改循环,审核方法必须使用 or_() 同时监听上游触发器和其自身的修改结果
@start()
def generate(self):
return "initial draft"
@human_feedback(
message="Approve or request changes?",
emit=["revise", "approved"],
llm="gpt-4o-mini",
default_outcome="approved",
)
@listen(or_("generate", "revise"))
def review(self):
return "content"
@listen("approved")
def publish(self):
return "published"
当结果为 "revise" 时,Flow 会路由回 review(因为它通过 or_() 监听 "revise")。当结果为 "approved" 时,Flow 继续执行 publish。这之所以有效,是因为 Flow 引擎免除了路由器的“仅触发一次”规则,允许它们在每次循环迭代时重新执行。
链式路由器
由一个路由器的结果触发的监听器本身也可以是一个路由器
@start()
def generate(self):
return "draft content"
@human_feedback(message="First review:", emit=["approved", "rejected"], llm="gpt-4o-mini")
@listen("generate")
def first_review(self):
return "draft content"
@human_feedback(message="Final review:", emit=["publish", "hold"], llm="gpt-4o-mini")
@listen("approved")
def final_review(self, prev):
return "final content"
@listen("publish")
def on_publish(self, prev):
return "published"
@listen("hold")
def on_hold(self, prev):
return "held for later"
@start() 方法运行一次:@start() 方法无法自循环。如果需要修改周期,请使用单独的 @start() 方法作为入口点,并将 @human_feedback 放在 @listen() 方法上。
- 同一方法上不能同时使用
@start() 和 @listen():这是 Flow 框架的约束。方法要么是起始点,要么是监听器,不能兼而有之。
最佳实践
1. 编写清晰的请求消息
message 参数是人员看到的内容。使其具有可操作性
# ✅ Good - clear and actionable
@human_feedback(message="Does this summary accurately capture the key points? Reply 'yes' or explain what's missing:")
# ❌ Bad - vague
@human_feedback(message="Review this:")
2. 选择有意义的结果
使用 emit 时,选择能自然映射到人员响应的结果
# ✅ Good - natural language outcomes
emit=["approved", "rejected", "needs_more_detail"]
# ❌ Bad - technical or unclear
emit=["state_1", "state_2", "state_3"]
3. 始终提供默认结果
使用 default_outcome 处理用户未输入直接按回车的情况
@human_feedback(
message="Approve? (press Enter to request revision)",
emit=["approved", "needs_revision"],
llm="gpt-4o-mini",
default_outcome="needs_revision", # Safe default
)
4. 使用反馈历史进行审计追踪
访问 human_feedback_history 以创建审计日志
@listen(final_step)
def create_audit_log(self):
log = []
for fb in self.human_feedback_history:
log.append({
"step": fb.method_name,
"outcome": fb.outcome,
"feedback": fb.feedback,
"timestamp": fb.timestamp.isoformat(),
})
return log
5. 处理路由与非路由反馈
设计 Flow 时,考虑是否需要路由
| 场景 | 使用 |
|---|
| 简单审核,仅需反馈文本 | 不使用 emit |
| 需要根据响应分支到不同路径 | 使用 emit |
| 带批准/拒绝/修改的审批门 | 使用 emit |
| 仅收集日志评论 | 不使用 emit |
异步人工反馈(非阻塞)
默认情况下,@human_feedback 会阻塞执行以等待控制台输入。对于生产应用,您可能需要异步/非阻塞反馈,以集成 Slack、电子邮件、Webhook 或 API 等外部系统。
提供程序抽象
使用 provider 参数指定自定义反馈收集策略
from crewai.flow import Flow, start, human_feedback, HumanFeedbackProvider, HumanFeedbackPending, PendingFeedbackContext
class WebhookProvider(HumanFeedbackProvider):
"""Provider that pauses flow and waits for webhook callback."""
def __init__(self, webhook_url: str):
self.webhook_url = webhook_url
def request_feedback(self, context: PendingFeedbackContext, flow: Flow) -> str:
# Notify external system (e.g., send Slack message, create ticket)
self.send_notification(context)
# Pause execution - framework handles persistence automatically
raise HumanFeedbackPending(
context=context,
callback_info={"webhook_url": f"{self.webhook_url}/{context.flow_id}"}
)
class ReviewFlow(Flow):
@start()
@human_feedback(
message="Review this content:",
emit=["approved", "rejected"],
llm="gpt-4o-mini",
provider=WebhookProvider("https://myapp.com/api"),
)
def generate_content(self):
return "AI-generated content..."
@listen("approved")
def publish(self, result):
return "Published!"
当引发 HumanFeedbackPending 时,Flow 框架自动持久化状态。您的提供程序只需通知外部系统并引发异常——无需手动持久化调用。
处理暂停的 Flow
使用异步提供程序时,kickoff() 返回 HumanFeedbackPending 对象,而不是引发异常
flow = ReviewFlow()
result = flow.kickoff()
if isinstance(result, HumanFeedbackPending):
# Flow is paused, state is automatically persisted
print(f"Waiting for feedback at: {result.callback_info['webhook_url']}")
print(f"Flow ID: {result.context.flow_id}")
else:
# Normal completion
print(f"Flow completed: {result}")
恢复暂停的 Flow
当反馈到达时(例如通过 Webhook),恢复 Flow
# Sync handler:
def handle_feedback_webhook(flow_id: str, feedback: str):
flow = ReviewFlow.from_pending(flow_id)
result = flow.resume(feedback)
return result
# Async handler (FastAPI, aiohttp, etc.):
async def handle_feedback_webhook(flow_id: str, feedback: str):
flow = ReviewFlow.from_pending(flow_id)
result = await flow.resume_async(feedback)
return result
关键类型
| 类型 | 描述 |
|---|
HumanFeedbackProvider | 自定义反馈提供程序的协议 |
PendingFeedbackContext | 包含恢复暂停 Flow 所需的所有信息 |
HumanFeedbackPending | 当 Flow 因反馈而暂停时由 kickoff() 返回 |
ConsoleProvider | 默认的阻塞控制台输入提供程序 |
PendingFeedbackContext
上下文包含恢复所需的一切
@dataclass
class PendingFeedbackContext:
flow_id: str # Unique identifier for this flow execution
flow_class: str # Fully qualified class name
method_name: str # Method that triggered feedback
method_output: Any # Output shown to the human
message: str # The request message
emit: list[str] | None # Possible outcomes for routing
default_outcome: str | None
metadata: dict # Custom metadata
llm: str | None # LLM for outcome collapsing
requested_at: datetime
异步 Flow 完整示例
from crewai.flow import (
Flow, start, listen, human_feedback,
HumanFeedbackProvider, HumanFeedbackPending, PendingFeedbackContext
)
class SlackNotificationProvider(HumanFeedbackProvider):
"""Provider that sends Slack notifications and pauses for async feedback."""
def __init__(self, channel: str):
self.channel = channel
def request_feedback(self, context: PendingFeedbackContext, flow: Flow) -> str:
# Send Slack notification (implement your own)
slack_thread_id = self.post_to_slack(
channel=self.channel,
message=f"Review needed:\n\n{context.method_output}\n\n{context.message}",
)
# Pause execution - framework handles persistence automatically
raise HumanFeedbackPending(
context=context,
callback_info={
"slack_channel": self.channel,
"thread_id": slack_thread_id,
}
)
class ContentPipeline(Flow):
@start()
@human_feedback(
message="Approve this content for publication?",
emit=["approved", "rejected"],
llm="gpt-4o-mini",
default_outcome="rejected",
provider=SlackNotificationProvider("#content-reviews"),
)
def generate_content(self):
return "AI-generated blog post content..."
@listen("approved")
def publish(self, result):
print(f"Publishing! Reviewer said: {result.feedback}")
return {"status": "published"}
@listen("rejected")
def archive(self, result):
print(f"Archived. Reason: {result.feedback}")
return {"status": "archived"}
# Starting the flow (will pause and wait for Slack response)
def start_content_pipeline():
flow = ContentPipeline()
result = flow.kickoff()
if isinstance(result, HumanFeedbackPending):
return {"status": "pending", "flow_id": result.context.flow_id}
return result
# Resuming when Slack webhook fires (sync handler)
def on_slack_feedback(flow_id: str, slack_message: str):
flow = ContentPipeline.from_pending(flow_id)
result = flow.resume(slack_message)
return result
# If your handler is async (FastAPI, aiohttp, Slack Bolt async, etc.)
async def on_slack_feedback_async(flow_id: str, slack_message: str):
flow = ContentPipeline.from_pending(flow_id)
result = await flow.resume_async(slack_message)
return result
如果您正在使用异步 Web 框架(FastAPI, aiohttp, Slack Bolt 异步模式),请使用 await flow.resume_async() 而不是 flow.resume()。在正在运行的事件循环内调用 resume() 会引发 RuntimeError。
异步反馈最佳实践
- 检查返回类型:暂停时
kickoff() 返回 HumanFeedbackPending——无需 try/except
- 使用正确的恢复方法:在同步代码中使用
resume(),在异步代码中使用 await resume_async()
- 存储回调信息:使用
callback_info 存储 Webhook URL、票据 ID 等
- 实现幂等性:为安全起见,您的恢复处理程序应该是幂等的
- 自动持久化:当引发
HumanFeedbackPending 时状态会自动保存,默认使用 SQLiteFlowPersistence
- 自定义持久化:如果需要,将自定义持久化实例传递给
from_pending()
从反馈中学习
learn=True 参数在人工审核员和记忆系统之间启用了反馈循环。启用后,系统通过从过去的人工修正中学习,逐步改进其输出。
工作原理
- 反馈后:LLM 从输出+反馈中提取可概括的经验,并以
source="hitl" 存储在记忆中。如果反馈只是批准(例如“看起来不错”),则不存储任何内容。
- 下次审核前:过去的 HITL 经验从记忆中被调取,并在人员看到输出前由 LLM 应用以改进输出。
随着时间的推移,人员会看到逐渐优化的预审核输出,因为每次修正都会告知未来的审核。
class ArticleReviewFlow(Flow):
@start()
def generate_article(self):
return self.crew.kickoff(inputs={"topic": "AI Safety"}).raw
@human_feedback(
message="Review this article draft:",
emit=["approved", "needs_revision"],
llm="gpt-4o-mini",
learn=True, # enable HITL learning
)
@listen(or_("generate_article", "needs_revision"))
def review_article(self):
return self.last_human_feedback.output if self.last_human_feedback else "article draft"
@listen("approved")
def publish(self):
print(f"Publishing: {self.last_human_feedback.output}")
第一次运行:人员看到原始输出并说“对于事实性声明,始终包含引文”。经验被提炼并存储在记忆中。 第二次运行:系统调取引文经验,预审核输出以添加引文,然后显示改进后的版本。人员的工作从“修复所有内容”转变为“捕捉系统遗漏的内容”。
| 参数 | 默认值 | 描述 |
|---|
learn | False | 启用 HITL 学习 |
learn_limit | 5 | 预审核时可回顾的最大历史经验数 |
关键设计决策
- 所有内容使用相同的 LLM:装饰器上的
llm 参数由结果缩减、经验提炼和预审核共享。无需配置多个模型。
- 结构化输出:在 LLM 支持的情况下,提炼和预审核都使用带有 Pydantic 模型的函数调用,否则退回到文本解析。
- 非阻塞存储:经验通过
remember_many() 存储,它在后台线程中运行——Flow 继续立即执行。
- 优雅降级:如果 LLM 在提炼过程中失败,则不存储任何内容。如果它在预审核过程中失败,则显示原始输出。两种失败都不会阻塞 Flow。
- 无需范围/类别:存储经验时,仅传递
source。编码管道会自动推断范围、类别和重要性。
learn=True 需要 Flow 具备可用记忆。Flow 默认自动获取记忆,但如果您通过 _skip_auto_memory 禁用了它,HITL 学习将被静默跳过。