我们如何保持长程智能体存活
构建一个持久的事件响应代理,它能在工作进程崩溃时存活、等待人类批准、安全地重试并恢复而不重复已完成的工作。
AI模型价格对比 | AI工具导航 | ONNX模型库 | Vibe Coding教程 | PLC在线仿真器 | Tripo 3D | Meshy AI | ElevenLabs | KlingAI | ArtSpace | Phot.AI | InVideo
一个事件响应代理已经收集了日志、关联了指标并提出了修复方案。它正在等待工程师批准更改。
然后工作进程Pod在部署期间被替换。
对话可能仍然存在于数据库中。提议的操作可能仍然在跟踪中可见。但有几个更难的问题仍然存在:
- 哪个步骤实际上已经完成?
- 修复请求是否已经发送?
- 失败的模型调用是否应该重试?
- 运行是否可以等待批准直到明天?
- 当工作进程返回时,执行应该从哪里恢复?
这些主要不是模型问题。它们是持久执行问题。
Temporal为我们提供了一种将编排编写为代码的方法,同时在崩溃、重启、超时和长时间等待中保留其执行状态。
它不会让代理更智能。它让代理的运行足够持久,能够在其周围的基础设施中存活。
1、核心思想一图概览
正常的代理框架通常帮助我们定义提示、工具、交接和推理循环。Temporal位于该逻辑之下,保持运行可恢复。
警报到达
↓
收集日志和指标
↓
要求模型进行诊断
↓
等待人类批准
↓
应用更改
↓
验证恢复
↓
关闭或升级事件
Temporal使用两个主要构建块对此建模:
Workflow(工作流)
确定性编排和持久状态
Activity(活动)
可能失败、超时或重试的外部工作
Workflow拥有顺序、分支、等待状态和关于下一步发生什么的决定。
Activity执行外部工作,例如调用LLM、查询监控API、写入数据库或重启服务。
这种分离很重要,因为Temporal可能会重放Workflow代码以重建其状态。因此Workflow代码必须是确定性的。网络调用、模型调用、数据库查询和现实世界的副作用应该放在Activity中。
2、Temporal实际持久化了什么
Temporal为每个Workflow执行记录仅追加的事件历史。历史捕获重要进度,如Workflow开始、Activity调度、Activity完成、信号、计时器、失败和重试。
当工作进程消失时,另一个工作进程可以重放该历史并重建Workflow的逻辑状态。已完成的Activity不会在重放期间简单地重复,因为它们记录的结果会返回给Workflow。
有一个重要的边界:
Temporal恢复Workflow的逻辑进度。它不会从进程停止的确切指令恢复任意Activity代码。
如果Activity在中途崩溃,该Activity可能会从头开始运行。这就是为什么重试安全的副作用仍然是我们的责任。
3、一个实用的事件修复工作流
以下示例保持核心机制可见。模型调用、监控调用和基础设施更改是Activity。Workflow通过信号协调它们并等待批准。
from dataclasses import dataclass
from datetime import timedelta
from temporalio import activity, workflow
from temporalio.common import RetryPolicy
@dataclass
class Incident:
incident_id: str
service: str
summary: str
@dataclass
class RemediationPlan:
action: str
reason: str
@activity.defn
async def collect_evidence(incident: Incident) -> dict:
# 在此调用日志和指标系统
return {
"service": incident.service,
"error_rate": 0.18,
"recent_deployment": True,
}
@activity.defn
async def propose_remediation(evidence: dict) -> RemediationPlan:
# 在此调用模型,提供证据和严格的输出模式
return RemediationPlan(
action="roll_back_latest_deployment",
reason="Error rate increased immediately after deployment",
)
@activity.defn
async def apply_remediation(plan: RemediationPlan) -> str:
info = activity.info()
idempotency_key = f"{info.workflow_run_id}-{info.activity_id}"
# 将密钥传递给部署服务,这样重试不会创建第二个回滚请求
return await deployment_api.apply(
action=plan.action,
idempotency_key=idempotency_key,
)
@activity.defn
async def verify_recovery(service: str) -> bool:
# 在此查询健康检查和监控
return True
@workflow.defn
class IncidentRemediationWorkflow:
def __init__(self) -> None:
self.approved = False
self.approved_by: str | None = None
@workflow.signal
def approve(self, approver: str) -> None:
self.approved = True
self.approved_by = approver
@workflow.run
async def run(self, incident: Incident) -> dict:
retry = RetryPolicy(
initial_interval=timedelta(seconds=2),
backoff_coefficient=2.0,
maximum_attempts=3,
)
evidence = await workflow.execute_activity(
collect_evidence,
incident,
start_to_close_timeout=timedelta(seconds=30),
retry_policy=retry,
)
plan = await workflow.execute_activity(
propose_remediation,
evidence,
start_to_close_timeout=timedelta(seconds=60),
retry_policy=retry,
)
await workflow.wait_condition(lambda: self.approved)
change_id = await workflow.execute_activity(
apply_remediation,
plan,
start_to_close_timeout=timedelta(minutes=2),
retry_policy=retry,
)
recovered = await workflow.execute_activity(
verify_recovery,
incident.service,
start_to_close_timeout=timedelta(seconds=30),
retry_policy=retry,
)
return {
"incident_id": incident.incident_id,
"approved_by": self.approved_by,
"change_id": change_id,
"recovered": recovered,
}
工作进程注册很小且故意无聊:
from temporalio.client import Client
from temporalio.worker import Worker
client = await Client.connect("localhost:7233")
worker = Worker(
client,
task_queue="incident-agent",
workflows=[IncidentRemediationWorkflow],
activities=[
collect_evidence,
propose_remediation,
apply_remediation,
verify_recovery,
],
)
await worker.run()
一个单独的客户端启动运行并稍后发送批准:
handle = await client.start_workflow(
IncidentRemediationWorkflow.run,
Incident("INC-1042", "checkout-api", "Error rate increased"),
id="incident-INC-1042",
task_queue="incident-agent",
)
await handle.signal(
IncidentRemediationWorkflow.approve,
"on-call-engineer",
)
result = await handle.result()
此示例故意比生产实现小,但责任边界是现实的:
- 编排和等待状态存在于Workflow中;
- 模型和工具调用存在于Activity中;
- 重试按Activity配置;
- 批准作为信号到达;
- 外部更改使用幂等键。
4、使价值显而易见的失败测试
可以使用Temporal开发服务器运行本地演示:
temporal server start-dev
启动工作进程和Workflow后,我们可以在collect_evidence完成但批准到达之前停止工作进程。
Workflow在Temporal中保持打开。事件历史已包含完成的Activity结果和当前等待状态。
当工作进程重新启动时,Workflow被重放,返回到批准等待,并在批准信号到达时继续。
这是保存代理消息和保留执行之间的实际区别。消息持久性告诉我们说了什么。持久执行告诉系统什么已完成、什么待处理以及接下来允许运行什么。
5、重试不意味着恰好一次
Temporal Activity使用至少一次执行模型。
考虑这种失败:
回滚请求到达部署API
↓
API接受请求
↓
工作进程丢失响应
↓
Temporal观察到超时
↓
Activity被重试
Temporal无法知道外部系统是否完成了第一个请求。如果没有幂等键或协调检查,重试可能会重复副作用。
这是最重要的生产经验之一:
Temporal使重试持久。我们仍然必须使副作用对重试安全。
对于支付、部署、配置或消息传递系统,这通常意味着幂等键、去重记录或写前读协调步骤。
6、人类批准成为正常的工作流状态
人类批准步骤不应该要求工作进程保持存活数小时。
Temporal信号允许外部客户端异步更改Workflow状态。Workflow可以使用workflow.wait_condition等待该状态。
这使长时间等待变得普通:
计划提议
↓
Workflow持久等待
↓
工作进程可能停止或重新部署
↓
批准信号稍后到达
↓
Workflow恢复
对于可以更改基础设施、发送客户通信、批准支付或修改生产数据的代理系统,这比将批准保持在一个长HTTP请求或内存代理会话中更清晰的控制边界。
7、权衡和Temporal未解决的问题
Temporal解决执行持久性,而不是语义正确性。
它不决定:
- 模型的诊断是否正确;
- 证据是否充分;
- 另一个推理回合是否值得其成本;
- 提议的操作是否被授权;
- 工具结果是否值得信任;
- 代理是否在没有进展的情况下重复自己。
我们仍然需要评估、策略检查、范围凭据、工具沙箱、可观察性和人类升级。Temporal可以可靠地保留和执行这些决策,但它不会替代它们。
它还引入了操作和设计权衡。Workflow代码必须保持确定性和与重放兼容。长时间运行的Workflow可能积累大量事件历史,因此非常活跃的代理循环可能需要继续作为新实例。
团队必须在运行仍然打开时仔细版本控制Workflow代码。Temporal服务器是MIT许可的,因此我们可以自托管它,但这增加了服务、持久性、升级、容量规划和值班所有权。
8、Temporal与LangGraph和代理框架
Temporal和代理框架解决系统的不同部分。
代理框架
提示、工具、交接、模型循环
图框架
状态转换、分支、连接
Temporal
持久执行、重试、计时器、等待、恢复
我们并不总是需要所有三个。一个简短的同步助手可能不需要Temporal的基础设施。一个长时间运行的编码、研究、事件、入职或批准代理可能立即受益于它。
9、我们会在哪里使用它
当代理运行满足以下条件时,Temporal非常合适:
- 持续时间超过一个请求或工作进程生命周期;
- 等待人类或外部事件;
- 跨越多个不可靠的服务;
- 执行不应该不必要地重复的昂贵步骤;
- 需要显式重试、截止日期和恢复;
- 必须在事件后可检查;
- 可能产生重要的现实世界副作用。
对于简单的聊天完成、简短的无状态工具调用或可以从头安全重启的后台作业,它可能是不必要的。
10、实际要点
生产代理不仅仅是模型加工具。它也是一个分布式过程,可以在网络失败、工作进程重启、批准需要数小时和外部系统返回模糊结果时保持活动。
Temporal为我们提供了一个持久的地方来保持该过程存活。
模型仍然可能推理不正确。工具仍然可能不安全。如果我们设计不当,重试仍然可能重复副作用。但一旦这些决策和边界被定义,Temporal可以确保运行不会在底层基础设施更改时简单地消失。
这就是真正的价值:
Temporal不会让我们的代理更智能。它让代理的执行足够持久以在现实世界中存活。
原文链接: How We Keep Long-Running AI Agents Alive with Temporal
汇智网翻译整理,转载请标明出处