发件箱与收件箱表:设计能够抵御故障的 Webhook
了解事务性出站队列、幂等入站队列、状态感知的延迟处理以及死信队列如何让 AWS、Azure 和 GCP 上的 Webhook 交付更加可靠。
Webhook看似是最简单的集成方式:一方发起HTTP POST请求,另一方负责处理。但实际上它存在分布式系统所有的风险,因为两个服务之间的网络可能会丢失请求、在传输过程中超时、重复发送相同的负载数据,或是打乱事件的顺序。如果将Webhook当作普通的CRUD调用来使用,最终会导致通知丢失、副作用被触发两次,进而让两个系统对实际发生的情况产生分歧。
本指南将详细介绍一种能在这些条件下正常运行的设计方案。您将了解为何简单的处理方式会失效,事务型消息队列如何确保发出的 Webhook 可靠,幂等性消息队列为何能让接收到的 Webhook 安全地重新发送,如何应对顺序错误的事件,如何隔离那些永远无法成功的请求数据,以及 AWS、Azure 和 Google Cloud 中哪些托管服务适用于各个环节。
为何常见的实现方式会导致数据丢失
假设有一个处理重要状态变化的 SaaS 后端,比如订单已完成配送或订阅已激活。合作伙伴会调用您的 API 来确认该操作,此时您的服务需要完成两项任务:
- 保存新的状态,例如将实体状态设置为
Active。 - 通过发送 Webhook 告知下游服务该实体已准备就绪。
这段直观的代码首先向数据库写入数据,然后在下一行发送HTTP请求。这属于双重写入:两个独立的系统依次被更新,彼此之间没有关联。
由此会直接产生两种故障模式:
- 在两个步骤之间进程崩溃。数据库仍认为该实体处于活跃状态,但实际上请求从未发送。你的记录是正确的,下游服务却毫不知情,直到有客户投诉才会被发现。
- 请求已发送,但事务失败。下游服务已收到该实体处于活跃状态的通知,但你的数据库已回滚操作,仍视其为失败。
无论哪种顺序都无法解决这个问题。将HTTP调用放在最前面会导致第二次失败;放在最后则会出现第一次失败。根本原因在于数据库提交与网络调用无法同时以原子方式执行,因此其间出现的任何崩溃或错误都会使双方数据不同步。
利用事务性消息队列可靠地发送Webhook
消息队列模式通过完全避免从请求路径发起HTTP调用,从而消除了双重写入的问题。相反,发送Webhook的意图会被转化为数据,而这些数据会与业务变更一起在同一个数据库事务中写入。要么两者都成功提交,要么都不提交。
消息队列流程的四个步骤
- 开启事务。业务操作会启动一个常规的数据库事务。
entities表(例如将状态设置为Active),并在outbox_events表中插入一行,其中包含下游服务应接收的完整数据。待发送表
下表为每个待处理通知存储一行数据。aggregate_type和aggregate_id用于标识该事件涉及的业务对象,event_type说明发生了什么,payload存储要传递的内容,而processed_at在中间节点确认已送达之前保持为空。查询processed_at为null的行即可得到中间节点需要处理的任务列表。注意,这里的行内注释使用单破折号;在PostgreSQL中注释需要双破折号(--),因此在执行语句前需进行修正。
CREATE TABLE outbox_events (
id UUID PRIMARY KEY,
aggregate_type VARCHAR(50), - e.g., 'Order' or 'User'
aggregate_id UUID, - e.g., Entity ID
event_type VARCHAR(100), - e.g., 'order.activated'
payload JSONB NOT NULL, - The exact webhook payload
created_at TIMESTAMP DEFAULT NOW(),
processed_at TIMESTAMP - Null until successfully sent
);
出站队列的保障与不足
如果在提交之后服务器崩溃,也不会有任何数据丢失:该行仍然存在于表中,中继会在下一次扫描时找到它。如果下游端点不可用,中继只需重新尝试,理想情况下采用指数退避策略,以避免给性能不佳的接收端带来过大压力。状态变更与通知意图再也不会出现不一致。
其代价是消息的交付会至少发生一次。中继节点可能在成功发送请求后,在更新processed_at之前就崩溃,这样一来相同的事件会在下一次运行时再次被发送。只有当接收方能够去重时这种情况才是可接受的,而另一端的收件箱模式正是为此设计的。在消息载荷或标头中包含出站消息的id,可以为接收方提供稳定的去重依据。如果运行多个中继实例,需确保没有两个工作进程能同时获取同一条记录;在PostgreSQL中,使用FOR UPDATE SKIP LOCKED来选择记录是实现这一目标的常见方法。
利用幂等收件箱安全接收webhook
现在将视角转向您的服务从合作伙伴或上游系统接收的webhook。
假设你的处理程序因执行复杂计算或等待其他服务持有的锁而耗时五秒。在你能响应之前,发送方的HTTP客户端可能会放弃,认为你从未收到该事件,于是再次发送。这样同一事件就会到达两次。如果你的处理程序每次运行都会发送邮件或创建记录,那么客户会收到两封邮件,而你也会得到一条重复的记录。
“收件箱模式”将接收webhook与对其进行处理分开。
收件箱流程的四个步骤
- 接收并验证。请求一到达,立即检查其HMAC签名,以确保它确实来自合作伙伴,且未被伪造或篡改。
webhook_inbox 表中,以合作伙伴独有的事件标识符作为键,并通过数据库唯一性约束进行保护。200 OK 状态。收件箱表
此处每一行记录了事件的发送者(partner_name)、该发送者的标识符(partner_event_id)、数据载荷、签名验证结果,以及处于 PENDING、PROCESSED 或 QUARANTINED 状态的 status。关键在于复合约束 UNIQUE(partner_name, partner_event_id):它能够将重复记录转化为无害的无操作处理。与发件箱表一样,单横线形式的注释需要改为 --, PostgreSQL才能接受该语句。
CREATE TABLE webhook_inbox (
id UUID PRIMARY KEY,
partner_name VARCHAR(50), - e.g., 'Stripe' or 'GitHub'
partner_event_id VARCHAR(100), - The unique ID from the sender
payload JSONB NOT NULL,
signature_verified BOOLEAN,
status VARCHAR(20), - 'PENDING', 'PROCESSED', 'QUARANTINED'
received_at TIMESTAMP DEFAULT NOW(),
processed_at TIMESTAMP,
UNIQUE(partner_name, partner_event_id) - Prevents duplicate inserts
);
为何该约束能发挥关键作用
由于处理程序仅负责验证、插入和返回操作,因此响应速度很快,发送方几乎不会出现超时情况。即便发送方尝试重试,哪怕连续重试十次,唯一性约束也会确保只有一次插入操作成功。处理程序应将由此产生的唯一性违反错误(或ON CONFLICT DO NOTHING结果)视为成功,并仍返回200 OK状态,否则发送方会不断重试已存在的事件。由于只存在一行数据,工作进程也只需执行一次相关操作。
有两点需要准确处理。首先,去重功能依赖于合作伙伴提供稳定的事件标识符;大多数Webhook服务提供商都会提供此类标识符,但仍需为每次集成进行确认。其次,工作进程在执行副作用操作后、标记该行已处理之前可能会崩溃,因此尽可能在同一个事务中完成业务变更和状态更新,并确保外部副作用操作也是可重试的。如需深入了解如何通过键值来去重请求,请参阅Node.js POST端点中的可重试键。
处理顺序错乱的事件
即便已控制住重复事件,也无法保证事件会按照产生的顺序到达。您的服务可能会先收到entity.completed,而后才收到entity.started。如果处理程序盲目地处理每个事件,就会试图直接将实体从draft状态转为completed状态,这要么会破坏其状态,要么会导致类似409 Conflict的错误。
根据状态机检查每次状态转换
解决方法是不再将事件视为修改状态的指令,而是将其视为需要验证的拟议状态转换。这有时被称为状态协调引擎,体现了事件源设计的理念:处理程序会将传入的事件与实体的当前状态进行比较,从而判断该转换是否合法。
下图展示了该决策过程。如果在实体仍处于草稿状态时出现了完成事件,由于前提条件尚未满足,函数会将该事件标记为延迟处理而非立即执行。注释中提到了两种处理延迟的方式:将该行保留在待办列表中稍后重试,或记录预期状态并等待缺失的事件。草稿实体上的启动事件属于有效转换,会被立即执行。可将其视为伪代码:return status: 'DEFERRED';并非有效的JavaScript语法,正确写法应为return { status: 'DEFERRED' };,而实际实现还需处理其他事件与状态的组合情况。
function processWebhookEvent(event, currentEntityState) {
if (event.type === 'entity.completed' && currentEntityState === 'draft') {
// The 'started' event hasn't arrived yet!
// We cannot transition from 'draft' directly to 'completed'.
// Option A: Leave it in the inbox and retry in 5 minutes.
// Option B: Store a "Projected State" and wait for the missing piece.
return status: 'DEFERRED';
}
if (event.type === 'entity.started' && currentEntityState === 'draft') {
return transitionTo('started');
}
}
作为自修复循环的延迟处理
假设存在这样一种订单:先收到“已发货”事件,而后才收到“已付款”事件。如果立即处理“已发货”事件,就会使订单进入模型不允许的状态。而使用具备状态感知功能的处理器后,处理顺序变为:
- 收到“已发货”事件后,评估器发现尚未付款,因此该事件会被延迟处理。
- 收到有效的“已付款”事件后,系统会更新订单状态。
- 随后重新尝试处理被延迟的“已发货”事件,此时发现其前置条件已满足,于是予以执行。
延迟处理的事件可以存放在专用的重试队列中,例如 Amazon SQS 或基于 Redis 的队列,然后由后台工作进程定期尝试重新处理这些事件。这样一来,流程能够拒绝无效的状态转换,同时最终会到达正确状态且不会丢失任何事件。不过需要为事件被延迟处理的时长设定限制:如果前置条件永远无法满足,该事件最终应被视为失败处理,而非无限重试,这引出了下一节的内容。
利用重试和死信队列隔离问题事件
有些事件无论重试多少次都永远无法成功:可能是有效载荷格式错误,或是引用的 ID 在数据库中并不存在。这类问题被称为毒药消息。如果处理程序过于简单,会无限次地重试这些消息,而一旦队列按顺序处理,一条错误消息就可能会阻塞其后所有正常的事件。
标准的防御措施是采用带有逐渐增加延迟的有限重试策略,之后再将无法处理的消息放入死信队列(DLQ)中。典型的处理流程如下:
- 第一次尝试失败,等待一分钟。
- 第二次尝试失败,等待五分钟。
- 第三次尝试失败,等待十五分钟。
- 第四次尝试仍然失败,则将该事件移至死信队列。
DLQ可以是您自己数据库中的表,也可以是托管队列的功能。关键在于后续的处理:DLQ中的事件应显示在内部管理界面中,并触发高优先级警报,因为每个事件都代表着系统无法处理的数据。工程师需要对此进行调查,修复映射错误或不良数据,然后重新播放该事件,使其通过正常的处理路径。应尽早实现这种重放功能;否则,要从DLQ中恢复数据就只能是在压力之下手动编辑数据库。
将设计映射到AWS、Azure和Google Cloud
发件箱和收件箱存在于您的关系型数据库中,但相关的组件(入口、队列、处理程序、DLQ等)则可以很好地映射到托管的云服务上,从而大大减轻运营负担。各提供商上的结构相同,只有产品名称会有所不同。
AWS
- 入口:Amazon API Gateway负责接收传入的webhook,请求到达后端之前会由Lambda授权器检查HMAC签名。
- 数据库:Amazon Aurora PostgreSQL存储业务表以及
webhook_inbox和outbox_events表,从而具备事务一致性保障。 - 队列与死信队列:SQS标准队列用于驱动异步处理,而配置好的SQS死信队列则会在消息数量超过最大接收限制时接收这些消息。标准队列本身为至少一次交付且不保持顺序,这也是上述幂等性和状态检查至关重要的原因之一。
processed_at 时间戳。Azure
- 入口:Azure API Management 负责接收 webhook 请求,验证签名后将请求转发至后端。
- 数据库:Azure Database for PostgreSQL Flexible Server 用于存储应用程序状态以及收件箱和发件箱表。
- 队列与死信队列:Azure Service Bus 负责消息的路由,具备内置的死信处理功能,即在达到预设的投递尝试次数后会自动将消息移至死信队列。
Google Cloud
- 入口点:Google Cloud API Gateway 负责处理传入的 HTTP Webhook 请求及身份验证。
- 数据库:Cloud SQL for PostgreSQL 用于存储关系型数据,包括各类表格。
- 队列与死信队列:Pub/Sub 用于异步传输消息。主订阅负责处理事件,而死信主题则会收集在规定的最大投递尝试次数后仍未被确认的消息。
有关服务连接的更多模式,例如 OAuth 和弹性 API 调用,请参阅六种用于可靠连接 Node.js 服务的集成模式。
关键要点
- 可靠的 webhook 系统是一个事件处理管道,而非一对 HTTP 接口。
- 切勿将更新数据库和调用远程服务视为两个互不相关的步骤;应在同一事务中写入发件箱记录,再由传递机制负责发送。