Webhook 交付系统设计实战:如何应对每天千万级事件

原文:https://dev.to/lovestaco/designing-a-webhook-delivery-system-for-10-million-events-a-day-2p5d(作者 @lovestaco)

团队里迟早会有人说出这句话,可能在规划会上,可能正盯着 Jira 上一条只写了三个字的工单。

“Webhook?那不就是一个 POST 请求。半天就能搞定。”

关于 POST 请求这部分,他们说得没错。

Webhook 真的就是这么简单。

你这边发生了什么事,你的客户想知道,你给他发一个 HTTP 请求。搞定。

然后你上线了,六周后,你在电话会议里向客户解释,为什么他们在自己服务器宕机期间错过了 4000 笔支付事件。

所以,我们来实际构建一下这个系统。

每天处理千万级事件,平均下来大约每秒 115 个,峰值时会高得多。

我会先构建一个最简陋的版本,然后故意把它搞崩,反复测试,直到我们得到一个能在真实客户环境下存活下来的设计。

版本 1:直接 POST

最直接的思路。事件在请求处理器中发生,你把它 POST 到客户的 URL,等待一个 200 OK 响应。

def on_payment_succeeded(payment):
    db.save(payment)
    requests.post(customer.webhook_url, json=payment.to_dict())  # 🙃
    return {"ok": True}


这在预发布环境里运行得非常完美,那里的“客户”是你开在另一个浏览器标签页里的 webhook.site。

在生产环境中它是这样的:

你的客户终端是一个 Rails 应用,跑在一个也运行着定时任务的小服务器上。

下午三点,定时任务启动,他们的服务器响应开始变得需要八秒钟,然后你的请求处理器就僵在那里,占用着线程,等待着别人的基础设施。

你的延迟图表出现尖峰。你的连接池被耗尽。

与你这些无关的自家用户开始看到超时。

然后是最糟的情况:你的请求超时了,进程继续前进,而那个事件……消失了。

你压根没把它记下来。它仅仅作为一个函数里的变量存在过,而该函数已经返回了。

Webhook,系统设计,消息队列,分布式系统,高并发

这里的 bug 不是“它很慢”。bug 是你让你的系统可用性取决于客户的可用性,而且你把它放在了关键路径上。

Webhook,系统设计,消息队列,分布式系统,高并发

版本 2:先记下来

分布式系统第一准则,坦白讲也是人生第一准则:在尝试执行之前,先把它记下来。

于是,处理器不再直接 POST。它先把事件写入一个数据库表,与触发它的业务变更在同一个事务里,然后返回。

这就是事务性发件箱模式。它巧妙地做了一件值得明说的事。

如果你先保存支付记录,再推送到一个队列,这是两个独立的系统,它们之间存在一个间隙。

如果在这个间隙中崩溃,你就得到了一笔没有事件的支付。

通过将事件行与支付记录放入同一个数据库事务,这两件事要么都发生,要么都不发生。没有间隙。

BEGIN;
  INSERT INTO payments (id, amount, status) VALUES (...);
  INSERT INTO webhook_outbox (customer_id, event_type, payload, status)
       VALUES (..., 'payment.succeeded', ..., 'pending');
COMMIT;


然后,由一个独立的工作线程轮询 pending 状态的行,并执行实际的 POST 操作。

你的处理器又变快了。你的事件也持久化了。

如果投递失败,它会失败在一个你能看见并重试的地方,而不是在一个已经消失的堆栈帧里。

但你只是用一个问题换来了另一个更隐蔽的问题。

你只有一个工作线程(或一个线程池),按顺序从一个队列中拉取任务。

客户 A 的终端超时需要 10 秒。

每个处理到客户 A 任务的线程都会挂起 10 秒钟。

与此同时,客户 B 到 Z 的事件就在队列里排在后面,明明可以投递,却哪儿也去不了。

这就是队头阻塞,是“吵闹邻居”问题穿上了队列的外衣。

一个终端有问题的客户拖慢了所有人。你最差的那个客户决定了所有客户的节奏。

Webhook,系统设计,消息队列,分布式系统,高并发

版本三:每位客户独享一条通道

解决方案是公平性,而公平性需要有人来执行。

在工作池前面放一个调度器。

它的全部职责是决定“下一个处理哪个事件”,并且不允许简单地选取最老的那个事件。

调度器会记录每个客户当前有多少 worker 正在处理其请求。

客户 A 已经有 3 个请求在处理中,而其并发上限就是 3?那就跳过它。取下一个客户的事件来处理。稍后再回来处理 A。

这就是按租户进行的并发限制,也是整个设计中杠杆效应最高的一点。

并发限制也是 Stripe 使用的核心速率限制原语,原因就在于此:它限制的是破坏范围,而不仅仅是计数请求。

效果是客户 A 的灾难现在有了上限。最多有 3 个 worker 卡在它那里。

工作池中的所有其他 worker 都可以愉快地服务其他所有客户。

现在,一个慢客户只会拖累自己,而痛苦落在正确的位置。

如果你想更进一步,调度器还可以实现加权公平性,这样你的企业级客户就不会被每小时发出一百万个事件的免费级客户饿死。

下面是路由逻辑,这实际上是系统的核心:

flowchart TD
    A[拉取下一个待处理事件] --> B{该客户是否已达到<br/>并发上限?}
    B -->|是| C[跳过,尝试下一个客户]
    B -->|否| D{目标端点熔断器<br/>是否打开?}
    D -->|是| E[暂停,直到冷却期结束]
    D -->|否| F[交给一个空闲的 worker]
    C --> A
    E --> A
    F --> G[POST 签名后的有效载荷]

    classDef decision fill:#f4d35e,stroke:#b8991f,color:#1a1a1a
    classDef start    fill:#e9ecef,stroke:#6c757d,color:#1a1a1a
    classDef action   fill:#5ee6c8,stroke:#1f9c86,color:#1a1a1a
    classDef wait     fill:#ff9a5c,stroke:#c26a33,color:#1a1a1a

    class B,D decision
    class A start
    class F,G action
    class C,E wait


一旦你实现了并发上限,那个熔断器分支就值得加上。

如果一个客户的端点连续失败了最近的 20 次尝试,你已经知道下一次也会失败。

别再浪费 worker 去验证了。

Webhook,系统设计,消息队列,分布式系统,高并发

版本四:正确发送一条 Webhook

现在让我们聚焦到最小的细节。一个 worker 接取了一个任务。接下来会发生什么?

设置较短的超时时间。十秒,而不是六十秒。一个响应缓慢的端点就是一个坏掉的端点,你不应该在查明情况时让它扣押一个 worker。

然后你查看返回的内容,重要的一步是:并非所有失败都是同一种失败。

  • 200201204:投递成功。标记它,然后继续处理下一个。
  • 连接被拒绝、超时、502503429:临时性故障。他们的服务器暂时有点问题。使用指数退避策略进行重试,并加上抖动。这样当他们的服务器恢复时,你不会在同一个毫秒内将所有重试积压的请求一股脑地打过去。AWS 写过一篇关于抖动的经典文章,值得花十分钟读一读。
  • 404410、DNS 无法解析、TLS 握手失败:永久性故障。URL 错误,或者端点已经不存在。在 24 小时内重试 12 次不是韧性,这只是你在生成无效流量,并延迟了客户发现其配置损坏的时机。快速且明确地让它失败。

第三种情况是团队常常忽略的,也正是它把你的重试队列变成了垃圾场。

顺带一提,有两件事是必须做的。

对有效载荷进行签名。 每个事件发出时,都在一个头部包含基于正文内容和时间戳计算的 HMAC。

你的客户使用共享密钥重新计算它,并确认事件确实来自你。

如果没有这个,你的 Webhook 端点就是一个任何人只要猜到 URL 就能向其发送虚假“支付成功”事件的地方。

将时间戳 包含在 签名内容中,这样被截获的请求下周就不能被重放给他们。

假设他们会处理两次。 重试意味着至少一次投递。

这不是你能通过工程手段消除的缺陷,而是问题本身的特性:如果你的请求超时了,你确实无法判断他们是否已经处理了它。

因此,给每个事件一个稳定的 id,告诉你的客户以此为键值,并清晰地写进文档。

“恰好一次投递”是营销话术。“至少一次投递”加上幂等性才是工程。

Webhook,系统设计,消息队列,分布式系统,高并发

版本5:死信队列是一项产品功能

重试总有用尽的时候。这种情况会发生。他们的端点在整个八小时重试窗口期间都处于宕机状态,已经没有更多可以尝试的了。

事件不会被删除,而是进入死信队列。

接下来这一点是我特别想强调的,因为大多数实现在这里就停了一步:一个没人能看到的死信队列,只是另一种更慢的数据丢失方式。

给它加个界面。在你的产品里做一个仪表盘,让客户能看到自己失败的投递记录。

事件类型、时间戳、尝试次数,以及从他们服务器实际收到的响应体。

最后那条信息能省掉大量工单,因为"下午3:04你们服务器返回了 502 Bad Gateway"这种具体信息,能立刻结束"webhook 坏了"这种争论——否则这种争论会拖上四天。

然后给他们一个重放按钮。

他们修好了部署,点击重放,事件重新进入 outbox 状态变为 pending,然后走完全相同的管道流程。

支持按时间范围批量重放,这样他们可以一键恢复整个故障窗口。

你刚刚把最糟糕的支持对话变成了一次自助操作。这是一个真正有价值的权衡。

Webhook,系统设计,消息队列,分布式系统,高并发

整体架构,端到端

flowchart LR
    APP[App writes event<br/>+ business change<br/>in one transaction] --> OUT[(Outbox)]
    OUT --> DISP[Dispatcher<br/>per-customer caps]
    DISP --> W[Worker pool]
    W -->|HMAC signed POST| CUST[Customer endpoint]
    CUST -->|2xx| DONE[Delivered]
    CUST -->|5xx / timeout| RETRY[Backoff + jitter]
    CUST -->|4xx permanent| DLQ[(Dead letter queue)]
    RETRY --> W
    RETRY -->|attempts exhausted| DLQ
    DLQ --> UI[Replay UI]
    UI -->|customer clicks replay| OUT

    classDef store  fill:#9d8cff,stroke:#5b4bcc,color:#1a1a1a
    classDef proc   fill:#5ee6c8,stroke:#1f9c86,color:#1a1a1a
    classDef ext    fill:#6ea8ff,stroke:#3565bd,color:#1a1a1a
    classDef bad    fill:#ff9a5c,stroke:#c26a33,color:#1a1a1a

    class OUT,DLQ store
    class APP,DISP,W,UI,DONE proc
    class CUST ext
    class RETRY bad


用一句话概括就是:发送之前先写下来,公平决定下一步发送给谁,用签名和短超时发送出去,对值得重试的失败进行重试,而对那些不值得重试的,让能真正修复问题的人看到它们。

Webhook,系统设计,消息队列,分布式系统,高并发

后续会咬你的那些坑

有几件事没法整齐地塞进版本迭代里,但它们一定会找上你。

顺序性。 有人会提这个要求。按客户排序投递意味着该客户的并发度为 1,也就是一个慢响应会阻塞他们的整个事件流。这是一个真实的权衡,不是免费功能。通常更好的做法是发送序列号让他们自己排序,或者发送"有东西变了,来拉取"的通知而不是状态本身。

载荷大小。 不要往 webhook 里塞 4MB 的对象。发送 ID 和事件类型,让他们调用你的 API 获取完整数据。轻量的载荷在 outbox 中存储成本更低,重试成本也更低,同时也避开了一个尴尬的问题:当状态在事件触发和他们读取之间发生了变化怎么办。

SSRF。 客户给你一个 URL,然后你的服务器去请求它。这是教科书式的服务器端请求伪造。解析主机名,拒绝私有范围和链路本地地址,并在重定向时重新检查,因为 http://customer.com/hook 重定向到 169.254.169.254 就是有人在尝试读取你的云元数据凭证。OWASP 有完整的列表

毒事件。 一个在反序列化时让 Worker 崩溃的事件会被无限重试,每次都让一个 Worker 挂掉。对你自己的失败也要限制尝试次数,不只是针对他们的。

Meme 想法 3 — 模板:他们不知道(派对上独自站在角落的人,内心独白)

  • 内心独白:"他们不知道我的 webhook 端点是个 Google 表格"

Webhook,系统设计,消息队列,分布式系统,高并发

所以,就半天功夫?

POST 请求部分确实只要半天。真的。

另外那 95% 的工作,是让事件保持存活的发件箱机制,是防止单个用户拖垮所有人的分发器,是能区分“可以重试”和“永远没戏”的重试分类器,以及能把数据丢失事故变成一个按钮的回放界面。

这些都不算什么黑科技。就是一张表,一个分发器循环,再加上一点对故障模式的严谨态度。

但正是这些,决定了你的 webhook 系统是客户愿意信任的,还是他们不得不围绕它编写防御性轮询代码的——而一旦他们错过一个事件,而你又无法告诉他们事件去向何处,他们立刻就会这么做。

先把日志写下来。其他一切都会从中衍生。

如果你一直在构建类似的 webhook 管道,我真的很想听听你目前卡在哪个版本。

我猜是第二版,而且导致问题的那个客户,你能直接说出名字。

原文:https://dev.to/lovestaco/designing-a-webhook-delivery-system-for-10-million-events-a-day-2p5d(作者 @lovestaco)

发布评论
全部评论(0)