Temporal Replay 2026 深度解析:让 AI Agent 拥有持久记忆,Serverless Workers 与 Workflow Streams 重新定义生产级 AI 工作流
引言:当「失败」从 bug 变成架构假设
分布式系统世界里,有一句老话:"网络是不稳定的,进程是会崩溃的,服务是会宕机的。" 传统的应对方式是把这些视为异常情况,用重试、用幂等、用补偿事务来打补丁。但 Temporal 的思路完全不同——它把这些统统视为正常情况,从架构层面就假设:任何操作都可能失败,任何进度都可能中断,任何状态都需要被持久化。
这就是 Durable Execution(持久化执行) 的核心哲学。
2026 年 5 月,Temporal 在 Replay 2026 大会上发布了一系列重磅更新:Serverless Workers 让 Temporal 可以跑在 AWS Lambda 上,Standalone Activities 让 Activity 不再必须依附 Workflow,Workflow Streams 让我们终于可以在持久化执行中看到 LLM 推理的实时 token 流,External Payload Storage 解决了 AI 场景下大模型输入输出的存储难题。这篇文章,我们就来深度拆解这些新特性,以及它们对生产级 AI Agent 工程实践的深远影响。
一、Durable Execution:重新理解「失败」这件事
1.1 传统可靠性的代价
在 Temporal 出现之前,工程师们为了保证分布式系统的可靠性,需要手工处理大量复杂逻辑:
重试逻辑:每个可能失败的外部调用都要写重试策略——指数退避、熔断、最大重试次数。这些逻辑散落在代码各处,一旦漏掉就埋下隐患。
幂等性:支付接口调用失败了,你敢不敢重试?不重试,钱没到账;重试了,可能扣两次钱。你得在数据库里维护一个 idempotency_key,在调用前查一下是否已经成功过。
补偿事务(Saga 模式):分布式事务不能用传统的两阶段提交,就得用 Saga——每个步骤都要写一个对应的补偿操作。一旦中途失败了,要按倒序执行所有补偿。一个涉及 5 个服务的 Saga,写出来的补偿代码可能比正向逻辑还长。
状态持久化:长流程(比如审批流、CI/CD 流水线)如果中途服务器重启,状态就丢了。你得自己存数据库、自己设计状态机、自己写恢复逻辑。
这些代码写出来,不仅复杂,而且脆弱——每个边界情况的处理都可能出错。
1.2 Temporal 的根本解法
Temporal 的思路是:把状态和执行流本身当作一等公民。
在一个 Temporal Workflow 里,你的代码是这样的:
@workflow.defn
class OrderWorkflow:
@workflow.run
async def run(self, order_id: str) -> str:
# 查询库存
inventory = await workflow.execute_activity(
check_inventory,
order_id,
start_to_close_timeout=timedelta(seconds=30),
)
if not inventory.available:
return "OUT_OF_STOCK"
# 扣款
payment = await workflow.execute_activity(
process_payment,
order_id,
start_to_close_timeout=timedelta(minutes=5),
)
if not payment.success:
# 取消库存预留
await workflow.execute_activity(
release_inventory,
order_id,
start_to_close_timeout=timedelta(seconds=30),
)
return "PAYMENT_FAILED"
# 发货
await workflow.execute_activity(
ship_order,
order_id,
start_to_close_timeout=timedelta(hours=2),
)
return "SUCCESS"
这段代码看起来平平无奇,但它的威力在于:整个 Workflow 的每次执行状态都被 Temporal Service 持久化了。
当你执行到 process_payment 这一步时,如果支付服务宕机了、Worker 进程崩溃了、甚至整个服务器断电了——Workflow 不会丢失任何状态。Temporal Service 会把 Workflow 恢复到 process_payment 调用之前的那个状态,然后重新执行(如果 process_payment 还没返回的话,会重新调用;如果已经返回了,就跳过)。
这就是 Temporal 的Replay 机制:Workflow 代码是确定性的(给定相同输入,总是产生相同输出),所以可以在任意时间点恢复并继续执行。
1.3 Workflow 与 Activity:两个抽象,两种命运
Temporal 把业务逻辑分成两层:
Workflow(工作流):编排逻辑。用 Temporal SDK 的特殊 API 写(workflow.execute_activity、workflow.sleep、workflow.wait、workflow.condition 等)。必须是确定性的——不能有随机数、不能直接读系统时间、不能调用外部 API。这些约束使得 Workflow 可以被安全地 Replay。
Activity(活动):实际干活的业务逻辑。普通的编程语言代码,可以调用任何外部服务(数据库、支付 API、LLM),可以有副作用,可以访问网络、时间、文件系统。Activity 有自动重试机制——默认会重试,而且重试策略可以精细控制。
这种分离带来了一个优雅的结果:业务逻辑(Activity)不需要知道「我在一个 Workflow 里」,它就是普普通通的函数。这使得 Temporal 特别容易与现有系统集成。
二、Replay 2026:五类新特性全面解读
Replay 2026 发布的内容可以分为五大类:开发效率提升、AI 场景强化、生态集成、生产运维、平台运营。我们逐一深入拆解。
2.1 Serverless Workers:Temporal 进了 AWS Lambda
这是本次发布中最具「范式转换」意味的特性。
传统上,运行 Temporal Worker 需要你自己管理计算资源:准备服务器或 K8s Pod、配置自动扩缩容、设计优雅关闭流程。Serverless Workers 则彻底改变了这个模式——你可以直接把 Temporal Worker 跑在 AWS Lambda 上。
# serverless_worker.py
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from workflows import YourWorkflow
from activities import your_activity
async def main():
client = await Client.connect("_NAMESPACE.account.tmprl.cloud:7233")
worker = Worker(
client,
task_queue="your-task-queue",
workflows=[YourWorkflow],
activities=[your_activity],
# 关键:新特性:允许 Serverless 模式
build_id_mode=WorkerBuildIdMode.AUTO_UPGRADE,
)
await worker.run()
# AWS Lambda handler
def handler(event, context):
return asyncio.run(main())
Serverless Workers 的核心价值:
- 零基础设施管理:不需要管理服务器、不需要配置自动扩缩容。Lambda 自动根据任务队列负载扩缩,最小可以缩到零。
- 成本优化:只有实际处理任务时才计费。对于脉冲式工作负载(突发流量),比常驻 Worker 便宜得多。
- 接入 Temporal Cloud 的自动管理:Temporal Cloud 现在可以自动调用、自动扩缩、自动优雅关闭 Lambda Worker。开发者只需要提供 Lambda 函数,Temporal 负责根据工作负载决定调用频率和并发度。
这对 AI 推理场景特别有意义:AI 推理请求量波动大,用 Lambda 处理峰值可以避免常驻 Worker 的资源浪费。
2.2 Standalone Activities:Activity 不再是 Workflow 的附属品
这是另一个架构层面的突破。
在 Temporal 传统模型里,Activity 必须作为 Workflow 的一部分来执行。也就是说,如果你想把一个普通的后台任务(定时任务、消息队列消费者、批量处理)纳入 Temporal 的可靠性保障,你必须把它包装在一个 Workflow 里。
Standalone Activities 解除了这个约束——Activity 现在可以独立运行:
# Standalone Activity 的用法
from temporalio.client import Client
from temporalio.worker import WorkflowWorker, ActivityWorker
async def main():
client = await Client.connect("...")
# 启动 Activity Worker(不再需要 Workflow Worker)
activity_worker = ActivityWorker(
client=client,
task_queue="standalone-task-queue",
activities=[process_batch_job, send_notification],
)
await activity_worker.run()
# 客户端独立调度一个 Activity
async def schedule_job(client: Client):
handle = await client.start_activity(
"process-batch",
args=[{"batch_id": "batch_001", "items": [...]},
task_queue="standalone-task-queue",
start_to_close_timeout=timedelta(hours=1),
id="job-batch-001", # 手动指定 ID,支持幂等
retry_policy=RetryPolicy(maximum_attempts=3),
)
result = await handle.result()
return result
Standalone Activities 的典型使用场景:
- 消息队列消费者:Kafka Consumer 处理消息时,经常需要把消息处理逻辑包装成幂等操作。Standalone Activity 让消费者逻辑天然具备 Temporal 的重试和可靠性保障,而不需要额外的 Workflow 编排。
- 定时任务:传统的 cron 任务丢了就丢了,没有状态记录。用 Standalone Activity 调度定时任务,处理过程会自动重试、完成状态可查、失败可告警。
- 批量处理管道:把批处理作业拆成独立的 Standalone Activity,可以单独重试、单独监控,不需要为了重试而设计整个 Workflow 结构。
Standalone Activities 的演进路径:当业务需求演进到需要多个 Activity 协同时,可以把相同的 Activity 代码直接嵌入 Workflow——不需要重写。这就是 Temporal 所说的「从简单开始,逐步演进」的开发体验。
2.3 Workflow Streams:持久化执行中的实时流
这是对 AI Agent 场景影响最大的新特性。
在 AI 应用中,用户最核心的体验需求之一是实时看到 LLM 的推理过程——token 一个一个吐出来,而不是等几秒后突然看到完整回复。传统 WebSocket/SSE 方案可以做到实时流,但它们没有持久化保障:一旦连接断开,服务重启,或者中间有网络抖动,LLM 的中间状态就丢失了。
Workflow Streams 把 Temporal 的 Signal & Update 机制与 LLM 实时流结合起来:
from temporalio.workflow import workflow_method, signal, update
import asyncio
@workflow.defn
class AIChatWorkflow:
@workflow.run
async def run(self, user_query: str) -> str:
# 初始化流处理器
stream_handler = workflow.streams.init_handler(
AIStreamHandler()
)
messages = [{"role": "user", "content": user_query}]
while True:
# 调用 LLM(作为 Activity)
response = await workflow.execute_activity(
call_llm_with_stream,
args=[messages, stream_handler],
start_to_close_timeout=timedelta(minutes=2),
)
if not response.done:
# 流式输出到客户端(通过 Workflow Streams)
for token_batch in response.tokens:
# Workflow Streams 会把这个 token 批量实时推送给监听者
await workflow.streams.emit(token_batch)
if not response.has_tool_calls:
return response.content
# 执行 tool call,继续循环
tool_result = await workflow.execute_activity(
run_tool_call,
args=[response.tool_calls[0]],
start_to_close_timeout=timedelta(minutes=5),
)
messages.append(response.to_message())
messages.append({"role": "tool", "content": tool_result})
# 客户端订阅 Workflow Streams
async def subscribe_stream(workflow_id: str):
async for token_batch in workflow.streams.subscribe(workflow_id):
print(token_batch, end="", flush=True)
Workflow Streams 的技术原理:Temporal 把 Workflow 的状态变更(通过 Signal/Update)暴露为流式接口。客户端可以实时订阅 Workflow 的中间状态变化,这些变化会被持久化——即使订阅者断开连接,重新订阅时也能从断点继续获取后续数据。
这解决了 AI Agent 开发中的一个核心矛盾:流式体验(实时响应)和可靠性(不丢数据)通常是对立的。Workflow Streams 让两者兼得。
2.4 External Payload Storage:大模型 I/O 的存储革命
这是最容易被人忽视,但影响最深远的更新。
LLM 的输入和输出可能非常大:一次 RAG 查询可能需要传入几千 token 的上下文;一次 AI Agent 的完整对话历史可能包含几十万 token;生成的代码文件、图片内容可能超过几 MB。
传统 Temporal 把这些数据存在 Temporal Server 的数据库里。对于大多数业务数据,这没有问题。但对于 AI 场景的大 payload,这会造成两个问题:
- 存储成本:Temporal Server 的数据库(PostgreSQL/Cassandra)存储大文件成本极高。
- 性能瓶颈:数据库不适合存储和传输超大文本块。
External Payload Storage 允许你把大 payload 存在外部存储(从 Amazon S3 开始支持,后续可以扩展到 GCS、Azure Blob 等),Workflow 里只存一个引用:
from temporalio.extensions import ExternalStorage
s3_storage = ExternalStorage(
provider="s3",
bucket="my-temporal-payloads",
region="us-east-1",
)
# 传入大文本时,External Storage 自动处理
large_context = load_rag_context(product_id)
# Temporal 会自动将 large_context 上传到 S3,只在 Workflow 状态里存一个 S3 引用
handle = await client.start_workflow(
RAGWorkflow.run,
args=[large_context, user_query],
id=f"rag-{product_id}-{user_id}",
task_queue="ai-queue",
storage=s3_storage, # 指定外部存储
)
result = await handle.result()
# 从 S3 取回结果
对于构建 AI Agent 应用来说,这个特性的意义在于:LLM 的输入输出不再受 Temporal Server 存储容量的限制。结合 Workflow Streams,可以构建真正的生产级 AI 数据处理管道。
2.5 生产级运维能力
Worker Versioning(GA):
传统上,升级 Worker 代码是个危险操作——正在执行的 Workflow 可能依赖旧版 Worker,新版 Worker 启动后可能破坏正在运行的任务。Worker Versioning 通过把运行中的 Workflow「钉」在启动它的 Worker 版本来解决这个问题:
worker = Worker(
client,
task_queue="production-queue",
workflows=[OrderWorkflow],
activities=[order_activities],
# 声明 Worker 版本
build_id="v2.3.1",
use_versioning=True,
)
新版 Worker(v2.3.1)部署后,新的 Workflow 会用新版运行,正在进行的 Workflow 继续用旧版直到完成。灰度发布、AB 测试都变得异常简单。
Task Queue Priority & Fairness(GA):
在多租户或多人协作的场景下,某个团队/租户的请求量暴增可能导致其他团队的请求被饿死。Task Queue Priority & Fairness 提供了租户级的流量控制和公平调度:
# temporal-config.yaml
taskQueue:
name: "production-queue"
priorityRules:
- priority: 10
filter:
workflowType: "PaymentWorkflow" # 支付流程最高优先级
- priority: 5
filter:
workflowType: "ReportWorkflow" # 报表次高
- priority: 1
filter: {} # 其他默认最低
fairnessPolicy:
enabled: true
weights:
team-a: 3 # 团队 A 权重更高
team-b: 2
team-c: 1
Nexus(GA for Python,Preview for TS/.NET):
跨 Namespace 的 Workflow 调用,传统的做法是每个团队维护自己的 Temporal Namespace,然后用 HTTP/gRPC 手动对接。Nexus 把跨 Namespace 的服务调用变成了一等公民:
# team-payments 服务(payments namespace)
from temporalio.nexus import service
@service
async def payment_service(ctx, amount: float, currency: str) -> PaymentResult:
result = await workflow.execute_activity(
process_payment,
args=[amount, currency],
start_to_close_timeout=timedelta(minutes=5),
)
return result
# team-orders 服务(orders namespace)直接调用
from temporalio.nexus import new_client
payment_client = new_client("payments") # 连接到 payments namespace
result = await payment_client.call_service(
"payment-service",
args=[100.0, "USD"],
)
Nexus 提供了跨 Namespace 的可观测性——你可以在 Temporal Web UI 里看到完整的跨服务调用链,延迟、错误率、重试次数一目了然。
三、AI Agent 场景:Durable Execution 的完美战场
3.1 为什么 AI Agent 最需要 Durable Execution
AI Agent 与传统软件有三个根本性的区别,让传统可靠性方案完全失效:
LLM 调用的不可靠性:LLM API 会超时、会限流、会返回错误响应。一次 Agent 循环可能调用 LLM 几十次,任何一次失败都可能导致整个 Agent 任务中断。
Tool Call 的副作用:Agent 调用外部工具(执行代码、查数据库、调第三方 API),这些操作有真实世界的副作用。失败后不能简单地重试,必须基于已有结果做决策。
中间状态的价值:LLM 的推理过程本身是有价值的中间状态——流式 token、tool call 历史、memory 快照。如果任务中断后这些状态全部丢失,Agent 就失去了「断点续传」的能力。
Temporal 天然解决了这些问题:
- LLM 调用作为 Activity:自动重试,内置指数退避,失败后 Workflow 从上一个成功节点恢复
- Tool Call 的幂等性保障:Activity 支持
idempotency_key,重复执行不会产生重复副作用 - 中间状态持久化:Workflow 的完整执行历史被保存,可以在任何时间点重新连接到并查看进度
3.2 与 Google ADK 和 OpenAI Agents SDK 的集成
Replay 2026 宣布了与两大主流 AI Agent 框架的深度集成:
Google ADK 集成:
# Google ADK + Temporal 集成
from google.adk.agents import Agent
from temporalio.contrib.adk import TemporalResumer
adk_agent = Agent(
model="gemini-2.0-flash",
tools=[database_query, file_search, api_call],
)
# TemporalResumer 让 ADK Agent 具备持久化能力
resumer = TemporalResumer(
workflow_id=f"agent-{session_id}",
task_queue="ai-agent-queue",
)
async def run_agent_session(user_input: str):
# Agent 的每次 LLM 调用和 tool 执行都作为 Temporal Activity
# 任何失败都会自动重试,从断点继续
result = await resumer.run(
agent=adk_agent,
user_message=user_input,
)
return result
OpenAI Agents SDK 集成(GA):
OpenAI Agents SDK 的 Sandbox(沙箱)功能与 Temporal 结合后,可以实现隔离环境中的持久化 Agent:
# OpenAI Agents SDK + Temporal 沙箱集成
from agents import Agent
from temporalio.contrib.openai_agents import TemporalSandboxRunner
agent = Agent(
model="o4-mini",
tools=[code_interpreter, file_tools],
)
runner = TemporalSandboxRunner(
sandbox_provider="e2b", # 隔离的云端沙箱
temporal_task_queue="agent-sandbox-queue",
)
# 每次 Agent 运行都是持久的——服务器宕机、网络抖动都不影响
async def run_isolated_agent(task: str):
result = await runner.run(agent, task)
# 即使 runner 进程崩溃,也可以通过 workflow_id 恢复
return result
3.3 实际案例:构建一个多步 AI 数据分析 Agent
让我们用 Temporal 实现一个完整的多步 AI 数据分析 Agent,感受 Durable Execution 在 AI 场景的威力:
import asyncio
from datetime import timedelta
from temporalio.client import Client
from temporalio.workflow import workflow_method, activity
# ============ Activities(实际执行逻辑)============
@activity.defn
async def fetch_data_source(source_id: str, query: str) -> dict:
"""从数据源获取原始数据"""
# 实际场景中这里调用数据库、数据湖、API
return {"rows": 10000, "columns": ["date", "revenue", "users"]}
@activity.defn
async def analyze_with_llm(raw_data: dict, objective: str) -> dict:
"""调用 LLM 进行数据分析"""
# 这里可以用 OpenAI、Claude、Gemini 等任意 LLM
prompt = f"分析以下数据,目标:{objective}\n数据:{raw_data}"
# 模拟 LLM 调用
return {
"insights": ["用户留存率下降 15%", "周末收入峰值明显"],
"recommendations": ["优化周末运营", "关注新用户转化"],
}
@activity.defn
async def generate_report(analysis: dict, format: str) -> str:
"""生成分析报告"""
# 生成 Markdown、PDF 或其他格式的报告
return f"# 数据分析报告\n\n## 发现\n{analysis['insights']}"
# ============ Workflow(编排逻辑)============
class DataAnalysisWorkflow:
@workflow.run
async def run(self, analysis_request: dict) -> dict:
request_id = analysis_request["id"]
objective = analysis_request["objective"]
sources = analysis_request["sources"]
results = []
# 第一阶段:并行获取多个数据源
fetch_tasks = [
workflow.execute_activity(
fetch_data_source,
source_id=src["id"],
query=src.get("query", ""),
start_to_close_timeout=timedelta(minutes=5),
retry_policy=RetryPolicy(maximum_attempts=3),
)
for src in sources
]
# 等待所有数据源完成
all_data = await asyncio.gather(*fetch_tasks)
# 第二阶段:LLM 分析
# 如果 LLM 调用超时/失败,Workflow 会自动重试
# 如果重试 3 次后仍失败,Workflow 可以 signal 给人工介入
analysis = await workflow.execute_activity(
analyze_with_llm,
args=[all_data, objective],
start_to_close_timeout=timedelta(minutes=3),
retry_policy=RetryPolicy(
maximum_attempts=5,
initial_interval=timedelta(seconds=5),
backoff_coefficient=2.0,
),
)
# 第三阶段:生成报告
report = await workflow.execute_activity(
generate_report,
args=[analysis, "markdown"],
start_to_close_timeout=timedelta(minutes=2),
)
return {
"request_id": request_id,
"status": "COMPLETED",
"report": report,
"insights": analysis["insights"],
}
# ============ 运行 ============
async def main():
client = await Client.connect("localhost:7233")
handle = await client.start_workflow(
DataAnalysisWorkflow.run,
args=[{
"id": "analysis-001",
"objective": "分析用户留存下降原因",
"sources": [
{"id": "mysql-users", "query": "SELECT * FROM user_events WHERE ..."},
{"id": "elasticsearch-logs", "query": "user_sessions"},
],
}],
id="data-analysis-001",
task_queue="data-analysis-queue",
)
# 即使这个进程退出,Workflow 也会继续执行
# 之后可以用 workflow_id 恢复并获取结果
result = await handle.result()
print(result)
这个实现的可靠性保障:
- 并行数据获取失败:任何数据源超时/失败,对应 Activity 自动重试,不影响其他数据源
- LLM 分析失败:最多重试 5 次,指数退避;超过重试上限后 Workflow 可以发出告警或等待人工介入
- 进程崩溃:Workflow 状态被持久化,重新启动客户端后可以继续获取结果
- 全程可观测:Temporal Web UI 展示每个 Activity 的执行时间、输入输出、错误信息、重试次数
四、为什么 Temporal 在 2026 年变得更重要
4.1 AI Agent 的「生产化」困境
2025 年是 AI Agent 的爆发年,但 2025 年底的调研显示,超过 70% 的 AI Agent 项目停留在 PoC(概念验证)阶段。主要原因是:Agent 的生产化比想象中难得多。
一个 AI Agent 要从 Demo 变成生产系统,需要解决:
- 可靠性:LLM 调用的不稳定性、外部工具的超时、AI 推理的不可预测性
- 状态管理:多轮对话的上下文管理、跨会话的长期记忆、agentic 循环的断点续传
- 可观测性:在 agentic loop 里发生了什么?哪一步 LLM 推理出了问题?Token 消耗如何跟踪?
- 成本控制:LLM 调用成本高昂,失败重试会成倍放大成本
Temporal 提供的 Durable Execution 基础设施,恰好覆盖了这些痛点:
| 需求 | Temporal 解法 |
|---|---|
| 可靠性 | Activity 自动重试 + Workflow Replay |
| 状态管理 | 完整执行历史持久化 + 中间状态可查 |
| 可观测性 | 内置 Web UI + OpenTelemetry + Billable Action Metrics |
| 成本控制 | Task Queue Priority + 细粒度 Activity 拆分 |
4.2 Temporal 的生态版图
Replay 2026 的一个重要主题是生态扩展:
- Temporal AI Partner Ecosystem:官方与多家 AI 供应商建立集成,包括 Google ADK、OpenAI Agents SDK、PydanticAI 等
- Rust SDK(Public Preview):Rust 开发者终于可以用第一方 SDK 构建 Temporal 应用了,这对于系统编程和高性能场景非常有意义
- AWS AI Competency:Temporal 获得了 AWS 在 Agentic AI 方向的官方能力认证,这意味着企业采购流程中 Temporal 更容易被接受
4.3 Temporal vs 其他方案
Temporal vs 消息队列(Kafka/RabbitMQ):
消息队列提供的是消息的可靠传递,但它不管理消息处理的状态。消费者崩溃后,消息可以重新投递给其他消费者,但处理进度(状态)就丢了。Temporal 则把「消息处理」和「状态管理」合二为一——不仅消息不丢,处理到哪一步也知道。
Temporal vs AWS Step Functions:
Step Functions 是托管的工作流服务,但它的状态机 DSL 表达能力有限,调试困难,成本随状态转换次数线性增长。Temporal 是开源的,状态机就是你的编程语言(Python/Go/TypeScript/Java),调试有完整的执行历史,成本模型基于 Workflow 运行时长而非转换次数。
Temporal vs 手工 Saga:
手工 Saga 的补偿逻辑是业务逻辑的一部分,容易出错,难以测试。Temporal 的 Saga 通过标准的 try/catch + Activity 重试机制实现,补偿逻辑是自动推导的(通过 Workflow Replay),不需要手工写反向操作。
五、从零搭建一个 Temporal 项目:完整实战指南
5.1 环境准备
# 安装 Python SDK
pip install temporalio
# 启动本地 Temporal Server(Docker 方式)
docker run -d \
--name temporal \
-p 7233:7233 \
temporalio/auto-setup:1.26.0 \
start \
--develop
# 验证服务启动
curl -s http://localhost:7233/health | jq .
5.2 项目结构
temporal_project/
├── pyproject.toml
├── workflows.py # Workflow 定义
├── activities.py # Activity 定义
└── worker.py # Worker 启动入口
5.3 一个完整的多步骤数据处理 Workflow
完整代码示例请参见第三章的实际案例。以下是项目运行方式:
# 启动 Worker
python worker.py
# 在另一个终端运行 Workflow
python run_workflow.py
5.4 Temporal Cloud 部署
对于生产环境,推荐使用 Temporal Cloud:
# 通过 Temporal Cloud CLI 认证
temporal auth login
# 查看 Namespace 信息
temporal namespace describe your-namespace
# 部署 Worker(连接 Temporal Cloud)
export TEMPORAL_HOST_URL="your-namespace.account.tmprl.cloud:7233"
export TEMPORAL_TLS_CERT=$(cat /path/to/ca.pem)
export TEMPORAL_TLS_KEY=$(cat /path/to/ca.key)
python worker.py # 自动连接到 Temporal Cloud
六、性能与架构:Temporal 是怎么做到的
6.1 Temporal Server 的存储模型
Temporal Server 使用 事件溯源(Event Sourcing) 模式:每个 Workflow 的状态不是直接存储的,而是存储导致状态变更的事件序列。
WorkflowExecutionStarted
→ ActivityTaskScheduled(check_inventory)
→ ActivityTaskCompleted(result=available)
→ ActivityTaskScheduled(process_payment)
→ ActivityTaskFailed(error=timeout)
→ ActivityTaskScheduled(process_payment) # 重试
→ ActivityTaskCompleted(result=success)
→ WorkflowExecutionCompleted
当需要恢复 Workflow 状态时,Temporal Server 从头 Replay 所有事件——这叫 Event Replay。对于复杂的 Workflow,这可能看起来很慢,但 Temporal 用 增量快照(Incremental Snapshot) 优化了这一点:每隔 N 个事件就拍一个快照,恢复时从最近的快照开始 Replay。
6.2 99.9999% 的可用性是怎么来的
Temporal Cloud 承诺 99.9999%(六个九) 的可用性,背后靠的是:
- Multi-region Replication(GA):Namespace 可以跨多个 AWS 区域复制,主区域故障时自动切换到备区域,RTO(Recovery Time Objective)≤ 20 分钟
- Multi-cloud Replication(GA):Namespace 可以跨 AWS 和 GCP 复制,满足多云合规要求
- 无状态 Worker:Worker 本身是无状态的,所有状态都在 Temporal Server。Worker 节点故障时,Workflow 自动调度到其他 Worker
- 历史服务分片:Temporal Server 用一致性哈希把 Workflow 分片到不同节点,单节点故障不影响其他 Workflow
6.3 成本模型解析
Temporal Cloud 的计费基于两个维度:
- Action Count:每次 Activity 执行的调用计为一个 Action
- Workflow Run Duration:Workflow 从启动到完成的时长
理解这个成本模型很重要:把大任务拆成多个 Activity 并行执行,通常比一个 Activity 做所有事情更快、更便宜。因为并行执行的 Activity 可以同时调度,总运行时间更短。
七、展望:Temporal 在 2026 年之后的路线
根据 Temporal 公开的路线图和技术趋势,以下几个方向值得关注:
更多语言的 First-party SDK:除了已经 GA 的 Go、Python、TypeScript、Java、.NET,以及 Preview 的 Rust,更多语言的 SDK(如 Ruby、PHP)可能进入官方支持列表。
更深入的 AI 框架集成:随着 AI Agent 框架生态的碎片化,Temporal 可能成为事实上的「AI Agent 可靠性层」——任何框架的 Agent 都可以通过 Temporal 获得生产级可靠性。
Workflow-as-a-Service:类似于 Vercel 对前端开发的变革,Temporal 可能推出更简化的托管服务,让开发者不需要理解 Server、Namespace、Cluster 等概念,直接上传 Workflow 代码即可。
更强的调试工具:Workflow 的 Replay 能力为调试提供了独特的可能性——可以在任意历史时间点「暂停」Workflow,检查状态,修改代码,然后继续执行。
总结:为什么每个 AI 应用开发者都应该了解 Temporal
分布式系统的可靠性问题,不是靠更多的测试、更严格的 code review、更完善的监控就能彻底解决的。你需要在架构层面假设一切都会失败,然后让系统在失败后能够自动恢复。
Temporal 做到了这一点。它不是又一个消息队列,不是又一个工作流引擎,而是把可靠性做成了编程模型的原生特性。你不需要写重试逻辑、不需要写补偿事务、不需要写状态恢复代码——你只需要写业务逻辑,Durable Execution 的基础设施会自动处理所有失败场景。
在 AI Agent 时代,这个能力变得尤为珍贵:LLM 调用的不可靠性、Tool Call 的副作用、多轮推理的状态管理——这些问题在 Temporal 面前都不再是难题。
如果你正在构建 AI 应用,或者你的分布式系统存在可靠性痛点,我强烈建议你花几个小时把 Temporal 的 Quick Start 过一遍。这可能是你今年学到的最有生产力的新技术。
Tags: Temporal|Durable Execution|AI Agent|Workflow Streams|Serverless Workers|分布式系统|可靠性工程|Python|Go|Replay 2026