编程 Temporal Replay 2026 深度解析:让 AI Agent 拥有持久记忆,Serverless Workers 与 Workflow Streams 重新定义生产级 AI 工作流

2026-07-23 06:44:59 +0800 CST views 7

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_activityworkflow.sleepworkflow.waitworkflow.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 的核心价值

  1. 零基础设施管理:不需要管理服务器、不需要配置自动扩缩容。Lambda 自动根据任务队列负载扩缩,最小可以缩到零。
  2. 成本优化:只有实际处理任务时才计费。对于脉冲式工作负载(突发流量),比常驻 Worker 便宜得多。
  3. 接入 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,这会造成两个问题:

  1. 存储成本:Temporal Server 的数据库(PostgreSQL/Cassandra)存储大文件成本极高。
  2. 性能瓶颈:数据库不适合存储和传输超大文本块。

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%(六个九) 的可用性,背后靠的是:

  1. Multi-region Replication(GA):Namespace 可以跨多个 AWS 区域复制,主区域故障时自动切换到备区域,RTO(Recovery Time Objective)≤ 20 分钟
  2. Multi-cloud Replication(GA):Namespace 可以跨 AWS 和 GCP 复制,满足多云合规要求
  3. 无状态 Worker:Worker 本身是无状态的,所有状态都在 Temporal Server。Worker 节点故障时,Workflow 自动调度到其他 Worker
  4. 历史服务分片:Temporal Server 用一致性哈希把 Workflow 分片到不同节点,单节点故障不影响其他 Workflow

6.3 成本模型解析

Temporal Cloud 的计费基于两个维度:

  • Action Count:每次 Activity 执行的调用计为一个 Action
  • Workflow Run Duration:Workflow 从启动到完成的时长

理解这个成本模型很重要:把大任务拆成多个 Activity 并行执行,通常比一个 Activity 做所有事情更快、更便宜。因为并行执行的 Activity 可以同时调度,总运行时间更短。


七、展望:Temporal 在 2026 年之后的路线

根据 Temporal 公开的路线图和技术趋势,以下几个方向值得关注:

  1. 更多语言的 First-party SDK:除了已经 GA 的 Go、Python、TypeScript、Java、.NET,以及 Preview 的 Rust,更多语言的 SDK(如 Ruby、PHP)可能进入官方支持列表。

  2. 更深入的 AI 框架集成:随着 AI Agent 框架生态的碎片化,Temporal 可能成为事实上的「AI Agent 可靠性层」——任何框架的 Agent 都可以通过 Temporal 获得生产级可靠性。

  3. Workflow-as-a-Service:类似于 Vercel 对前端开发的变革,Temporal 可能推出更简化的托管服务,让开发者不需要理解 Server、Namespace、Cluster 等概念,直接上传 Workflow 代码即可。

  4. 更强的调试工具: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

推荐文章

JavaScript中的常用浏览器API
2024-11-18 23:23:16 +0800 CST
Elasticsearch 聚合和分析
2024-11-19 06:44:08 +0800 CST
Web 端 Office 文件预览工具库
2024-11-18 22:19:16 +0800 CST
使用 `nohup` 命令的概述及案例
2024-11-18 08:18:36 +0800 CST
初学者的 Rust Web 开发指南
2024-11-18 10:51:35 +0800 CST
Go 并发利器 WaitGroup
2024-11-19 02:51:18 +0800 CST
程序员茄子在线接单