用 Temporal 独立 Activity 编排耐重试的 Webhook 队列

订单创建、支付完成或用户注册后,应用往往要向另一个团队提供的 URL 发送 Webhook。困难不在那一次 HTTP POST,而在网络中断、接收端返回 500、Worker 在发送后崩溃,或者上游把同一事件提交两次时,系统应该怎样继续。

Temporal 的 Standalone Activity 把一个普通 Activity 直接作为持久任务提交,不必先写 Workflow。平台记录任务、调度执行、按策略重试,并提供可查询的任务标识。业务数据仍应保存在应用自己的数据库;Temporal 服务也仍需要由团队或 Temporal Cloud 运维。本文沿着原教程六个模块,分别解释提交、幂等、去重、限速、心跳恢复和 Workflow 复用。

作者与来源:Nikolay Advolodkin,Temporal Staff Developer Advocate。原文 Build a Job Queue with Standalone Activities 标示更新于 2026-06-15;本文于 2026-10-05 读取完整正文,并核对官方课程仓库提交 42f50152433037262b9baf1e42cb66cb72da5284 中的六个 Python 模块、接收端与依赖声明。本文为授权翻译与技术整理;安全修正明确标注,所有代码仅静态审核。

Temporal Webhook 流程图:调用方用稳定事件ID提交,Temporal处理运行中冲突和完成后重用,Worker投递时携带接收端幂等键,完成副作用后再推进心跳检查点。
未完纪原创技术示意图。平台调度去重、接收端幂等和批次恢复分别解决不同问题。

先确认实验环境与版本

原教程推荐免费的 Instruqt 浏览器实验,沙箱会启动 Temporal 服务、Web UI 和 Webhook 接收端。页面公开可读,实验入口是否需要注册以平台为准;托管 Temporal 或外部 Webhook 服务的费用,不应与教程“免费实验”混为一谈。

本地路径需要 Git、Python 3.11 或更新版本、uv,以及支持 Standalone Activities 的近期 Temporal CLI。所读取的 pyproject.toml 声明 temporalio>=1.27 和 httpx>=0.27;这不是严格锁定版本。仓库 README 记录的是作者曾核验 Python SDK 1.27.2 和 v1.6.2-standalone-activity CLI 预发布版本的情况,不是本文重新验证的结果。

原文的本地准备步骤如下,本文没有执行:

git clone https://github.com/temporalio/edu-standalone-activities.git
cd edu-standalone-activities/python/course-repo
uv sync

课程按 exercise/01-... 至 exercise/06-... 提供带 TODO 的起始代码,并在对应 solution/ 目录提供答案。原文建议分别启动接收端、Temporal dev server 和模块 Worker,再提交任务:

# 在 course-repo 目录:只在隔离本地实验中使用
uv run python server/webhook_receiver.py

# 在另一个终端启动开发服务
temporal server start-dev --ui-port 8233

# 在模块目录启动 Worker
cd exercise/01-durable-job-queue
uv run python -m webhooks.worker

# 在另一终端的相同模块目录提交
uv run python -m webhooks.send_standalone evt_001
curl http://localhost:9000/_received

本地启动限制:原文明确把本地运行列为不提供支持的路径。本文静态查看的代码位于各模块的 src/webhooks/,根项目构建配置又指向 src/webhooks;实际安装与导入路径是否已由沙箱环境配置,需要在你选择的版本里核实。上述命令是原教程的启动路径,不是本文实测成功的安装方案。

接收端必须先收窄暴露范围:仓库接收端实际监听 0.0.0.0:9000,并不只监听 localhost。若在本地演练,至少应把创建服务器的地址改为下面的回环地址,并阻止实验端口对外开放。这是本文编辑修正,未执行验证:

server = HTTPServer(("127.0.0.1", 9000), EchoHandler)

模块一:直接提交一个持久任务

Activity 仍由 @activity.defn 标记。它是不是“独立 Activity”,取决于调用方式,而不是另一个特殊装饰器。原始发送函数只做一件事:向 URL POST JSON,遇到错误状态时抛出异常,成功时返回状态码。

import httpx
from temporalio import activity
from .shared import WebhookDelivery

@activity.defn
def deliver_webhook(req: WebhookDelivery) -> int:
    response = httpx.post(req.url, json=req.payload, timeout=5.0)
    response.raise_for_status()
    return response.status_code

调用方使用 client.execute_activity 提交并等待结果;使用 client.start_activity 则可以先取得句柄,稍后读取结果。下面是原调用关系,依赖课程模块中的共享类型和任务队列常量:

from datetime import timedelta
from temporalio.client import Client
from .shared import TASK_QUEUE, WEBHOOK_RECEIVER_URL, WebhookDelivery

client = await Client.connect("localhost:7233")
result = await client.execute_activity(
    deliver_webhook,
    args=[WebhookDelivery(
        url=WEBHOOK_RECEIVER_URL,
        payload={"event_id": "evt_001", "type": "order.created"},
        event_id="evt_001",
    )],
    id="deliver-evt_001",
    task_queue=TASK_QUEUE,
    start_to_close_timeout=timedelta(seconds=10),
)

这段带 await 的代码应放在异步函数或支持顶层 await 的环境中,不是可直接粘贴到普通脚本顶层的完整程序。课程的完整文件通过 asyncio.run(main(...)) 启动;同步 HTTP Activity 由 Worker 的 ThreadPoolExecutor 执行,Worker 注册 Activity 并持续轮询相同 Task Queue。

id 给任务稳定地址,便于查询、取消或终止。start_to_close_timeout 限制单次执行的时长,不等于包括所有重试在内的业务总期限。原教程举例:第一次遇到临时 503 时,默认重试会先等待初始间隔,再指数退避,任务仍显示运行中,attempt 增加。永久错误不会因为重复发送自动变成成功,所以生产中应明确重试次数、期限和不可重试错误。

模块二:让副作用经得起重试

HTTP 请求已经到达接收端,但 Worker 没收到响应,或者收到响应后尚未报告成功就崩溃,都会留下一个不确定窗口。Temporal 可以再次执行整个 Activity 函数,因此必须按可能多次执行设计外部写入。有限重试、取消或永久故障也可能最终失败,不能把“至少一次执行语义”理解为无限条件下保证业务成功。

课程在第二模块添加稳定的幂等键:

headers = {"Idempotency-Key": f"webhook:{req.event_id}"}
response = httpx.post(
    req.url, json=req.payload, headers=headers, timeout=10.0
)
response.raise_for_status()

键取自逻辑事件 ID,所以同一事件的每次尝试携带同一个值。每次执行都生成 uuid.uuid4() 或随机折扣码,会让接收端认为它们是不同事件,无法消除重复。需要随机生成的业务值,应由调用方生成一次并作为稳定输入传入。

幂等是接收端协议,不是加了一个请求头就天然成立。接收端需要持久保存键、已处理结果和必要的载荷指纹,让“记录幂等键”和“提交业务副作用”保持一致;相同键却不同载荷应有明确拒绝或冲突规则。键还应区分租户与业务事件范围,避免不同客户互相覆盖。

第二模块的答案会故意在前两次成功 POST 之后抛出 ApplicationError,用于观察重复请求是否被接收端合并;提交端设置 RetryPolicy(maximum_attempts=5)。这属于故障演示逻辑,发布业务代码时不能保留这个人为失败分支。本文没有执行故障注入,也没有声称观察到重试结果。

模块三:处理调用方的重复提交

接收端幂等解决的是 Activity 重试导致的重复副作用。另一类问题是上游对同一事件调用两次 start_activity。运行中已有同 ID 任务时,默认冲突会抛出 ActivityAlreadyStartedError;课程使用以下策略返回现有任务句柄:

from temporalio.common import ActivityIDConflictPolicy

handle = await client.start_activity(
    deliver_webhook,
    args=[req],
    id=f"deliver-{req.event_id}",
    task_queue=TASK_QUEUE,
    start_to_close_timeout=timedelta(seconds=30),
    id_conflict_policy=ActivityIDConflictPolicy.USE_EXISTING,
)

此处 req 是调用方已经构造的 WebhookDelivery。两次提交若发生在同一任务仍运行时,取得的句柄指向相同执行;重复提交不会变成第二个 Worker 任务。课程答案还特意暂停一秒,让第二次调用较容易落在这个窗口里;这只是教学安排。

完成以后是另一套规则。原文问答明确指出:如果第一次已经完成,之后再提交相同 ID,默认 ALLOW_DUPLICATE 重用策略可以启动新执行。需要同时限制完成后的重复时,还要考虑 id_reuse_policy=ActivityIDReusePolicy.REJECT_DUPLICATE,并核对当前 SDK、服务版本及保留策略。平台任务历史不是业务数据库永久去重表,接收端幂等仍不可省略。

模块四:限速与优先级

下游每秒只能接收一定数量的请求,而 Worker 按自己的能力尽快发送时,持续 429 会把重试变成额外流量。重试可以处理短暂失败,却不能单独修复总体发送速率超过接收能力的问题。

课程在 Worker 上设置:

worker = Worker(
    client,
    task_queue=TASK_QUEUE,
    activities=[deliver_webhook],
    activity_executor=executor,
    max_activities_per_second=2.0,
)

executor 是已创建的线程池。这一配置限制这个 Worker启动 Activity 的速率;不是线程池大小,也不是整个下游服务每秒请求数的保证。两个相同 Worker 的限额会相加;若一个 Activity 内发送一批请求,Activity 启动速率也不能直接等同于 HTTP 请求速率。原文建议需要跨 Worker 控制队列速率时检查 max_task_queue_activities_per_second。

排队任务由 Temporal 保存,但配置了超时、取消或有限重试的任务仍可能结束为失败。应同时观察队列等待时间、下游 429 和重试量,而不只是看单个 Worker 是否忙碌。

原文还介绍 Priority(priority_key, fairness_key, fairness_weight):较小的 priority key 用来把紧急任务安排到积压任务前面,公平性字段用于控制不同租户的份额。它依赖相应服务和 SDK 功能;不能把优先级当作严格实时期限,更不能用它代替接收端容量规划。

模块五:用心跳保存批次进度

一个 Activity 若要处理几十个 Webhook,失败后全部重做可能很浪费。课程把已经完成的数量写入 heartbeat,下一次尝试通过 activity.info().heartbeat_details 读取检查点,再从对应索引继续:

info = activity.info()
start_index = 0
if info.heartbeat_details:
    start_index = info.heartbeat_details[0]

delivered = start_index
for i in range(start_index, len(req.items)):
    item = req.items[i]
    response = httpx.post(req.url, json=item, timeout=5.0)
    response.raise_for_status()
    delivered += 1
    activity.heartbeat(delivered)

这是原批处理逻辑的核心。提交端示例设置单次执行五分钟、heartbeat_timeout 五秒;课程为了便于演示,还在每项后等待一秒。缺失心跳达到超时时,服务可以判定这次尝试失去活性,并按重试策略重新派发,而不用一直等单次执行的总超时。心跳也是运行中 Activity 接收取消信号的一条路径。

检查点没有消除重复窗口。所读取的第五模块批量发送代码没有携带 Idempotency-Key。POST 成功后到 heartbeat 被服务接收、持久化之间仍可能崩溃;下一次尝试会重发已产生副作用的项。SDK 还可能节流心跳发送,因此不能假设每次本地 activity.heartbeat() 调用都已经独立持久化。心跳超时也不代表旧 HTTP 请求已被硬性停止,旧尝试与新尝试可能存在副作用重叠。

下面是对上述批次片段的编辑修正:每项使用调用方提供的稳定事件 ID,并仍在 HTTP 成功后推进检查点。这里只展示局部变化,依赖接收端正确实现持久幂等,未做执行测试:

item = req.items[i]
headers = {"Idempotency-Key": f"webhook:{item['event_id']}"}
response = httpx.post(
    req.url, json=item, headers=headers, timeout=5.0
)
response.raise_for_status()
delivered = i + 1
activity.heartbeat(delivered)

生产事件 ID 应在业务范围内稳定且唯一。课程生成的 item_000 等名称只适合演示,若每批都从同名 ID 开始,直接拿来做跨批去重会把不同事件误当作同一件事。恢复时还应校验检查点类型和范围,并保持批次内容与顺序稳定。

检查点必须在确认副作用完成之后推进;提前标记完成会在崩溃后跳过尚未投递的项。即便顺序正确,仍应接受“最后若干项可能重试”,通过幂等让这种重试安全。本文没有运行仓库中的终止 Worker、重置接收端等演示脚本。

模块六:把相同 Activity 放进 Workflow

当任务从“发一个请求”扩展到多步业务流程时,可以保留 Activity 函数,把它作为 Workflow 中的一步。原教程的包装层很薄:

from datetime import timedelta
from temporalio import workflow

with workflow.unsafe.imports_passed_through():
    from .activities import deliver_webhook
    from .shared import WebhookDelivery

@workflow.defn
class WebhookWorkflow:
    @workflow.run
    async def run(self, req: WebhookDelivery) -> int:
        return await workflow.execute_activity(
            deliver_webhook,
            req,
            start_to_close_timeout=timedelta(seconds=10),
        )

HTTP 请求仍在 Activity 内执行;Workflow 调用 Activity,并把重试、超时与可见性纳入编排。imports_passed_through 让指定模块通过 SDK 的 Workflow 导入处理机制,它不是允许把任意不确定的网络操作放进 Workflow 的许可。

“同一个函数可以复用”不等于“每个教学模块自动继承了前面所有加固”。本次读取的第六模块发送函数回到了没有幂等键的基础版本。因此,实际迁移时应复用经过幂等、URL 校验和重试分类处理的 Activity,而不是把课程最简版本当作完整生产实现。

原文以 Coinbase 在 Replay 2026 分享的迁移为背景,提到每日 2 亿至 6 亿任务、186 个 namespace 的规模。这是原作者援引的案例陈述,本文未独立核验演讲数据,不能据此预测本实验或你的系统吞吐。

源码中需要实质补齐的安全边界

URL 与 SSRF:原 Activity 直接使用 httpx.post(req.url, ...)。课程常量指向本地接收端,但如果业务把这个字段开放给不可信调用方,Worker 就可能被诱导访问内网或元数据服务。实际系统应限制允许的协议、主机和端口,约束重定向并在网络层控制外连;解析一次 URL 或一次 DNS 结果并不能独自解决所有地址变化问题。日志也不应完整记录带敏感查询参数的 URL。

错误分类:原函数对错误状态统一抛异常,若没有策略,永久 4xx 也可能反复重试。可按业务协议把不可恢复的请求错误标为不可重试,并保留 408、429、部分 5xx 等需要重试的情况。401 或 403 是否能通过凭据刷新恢复,必须由实际协议决定。第二模块中 HTTP 超时与单次 Activity 超时都为十秒,缺少额外调度、记录和返回的裕量;应让各层期限互相协调。

教学接收端不是生产服务:它监听所有网卡;/_received 会返回已收到的完整载荷,/_reset 会清空记录,/_rate_limit 会改变限速,均未鉴权。幂等信息保存在进程内存列表里,重启或重置就丢失;收到的载荷没有明确大小上限,记录也会不断累积。它适合封闭沙箱内观察行为,不适合承载真实客户 Webhook。回环监听只缩小网络暴露,不能替代生产鉴权、限额和持久化设计。

数据与凭据:Activity 输入和心跳内容会交给 Temporal 服务,接收端也会保存载荷。不要传入不必要的秘密或个人信息;敏感数据的加密、脱敏和保留期限应在调用前规划。本文没有发现这些示例中嵌入真实生产令牌,但也没有审计整个仓库及所有依赖。

审查结论:六模块的教学主线成立,但耐重试的副作用需要平台去重、接收端持久幂等和正确检查点共同配合。本文没有启动服务、发送 Webhook、制造崩溃、调用收费 API 或测试吞吐;静态审查与文件核对不构成执行测试,也不保证不存在其他漏洞。

归属与许可:原作者 Nikolay Advolodkin / Temporal;本文保留来源与作者归属。所读取课程仓库树未见可确认的根 LICENSE,未擅自指定开源转载许可;正文、代码与配图经授权使用。图示、回环监听和批次幂等修正为本次编辑新增。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容