SNS、SQS 与异步任务的可靠处理

# SNS、SQS 与异步任务的可靠处理

入队成功只表示工作已被接收,业务成功还要等消费者处理、持久化并能回读结果。这是 AI 长任务、Web3 事件通知和性能数据清洗共同面对的问题。

服务基本概念复用 SQS 原笔记,事务与消息协作复用 Aladdin 的最终一致性与对账。

# SNS、SQS、EventBridge 怎么分工

服务 主要作用 一个例子
SNS 发布订阅,将同一事件分发给多个订阅者 “报告已生成” 同时通知校验与通知服务
SQS 缓冲待处理工作,让消费者按能力领取 模型任务、事件验证、数据清洗
EventBridge 发生指定事件或到达指定时间时,触发对应服务 订单支付后通知下游;每小时触发评分更新

原项目的一条链路是:Node 生成结果 → SNS → SQS → verifier Lambda → 回读公开 API。SNS 返回成功,只证明发布被接受,还需检查投递、队列、消费者和回读结果。

# EventBridge 能做什么,为什么要用

EventBridge (opens new window) 是 AWS 的事件与调度服务。它负责在指定时间,或某件事发生后,自动通知对应的程序去处理。常见两种用法:

  • 定时触发:每小时通知后端更新 Agent 评分,不需要用户打开页面才执行。
  • 事件触发:收到 “订单已支付” 事件后,按规则通知发货、积分等服务;订单服务只发布事件,不必逐个调用下游接口。

为什么用它:把触发时间、事件匹配、投递和失败重试交给 AWS 管理,减少自建调度与通知程序的维护工作。它不是必需品:简单定时任务也能用 cron,只是要自行负责程序运行、失败处理和监控。

cron 是 Linux/Unix 中常用的定时任务工具,可以按设定时间自动执行命令或脚本,例如每小时调用一次评分更新接口;它通常运行在自己管理的服务器上。

以 Aladdin 的评分更新为例:每小时到点 → EventBridge 调用评分快照接口 → 后端 Worker 更新评分。真正计算评分的是 Worker,EventBridge 只负责触发调用。通过 API Destination (opens new window) 可以直接调用已有 HTTPS 接口,不必额外写一个只负责转发请求的 Lambda。

配置里的几个名称分别做什么

  • Event Bus(事件总线):接收 “订单已支付” 这类事件。
  • Rule(规则):指定哪些事件交给哪些目标;已有定时 Rule 也可以按时间触发。专门的定时调度还可使用 EventBridge Scheduler,不必先发送业务事件。
  • API Destination(接口目标):配置要调用的 HTTPS 地址、请求方法和调用速率。
  • Connection(连接配置):管理调用接口所需的 API Key、Basic 或 OAuth 凭据;目标接口仍要验证身份与权限。

调用成功不等于业务完成。API Destination 等待接口响应的超时为 5 秒,耗时任务应先可靠保存或入队,再异步执行。重试可能重复调用,后端需要幂等处理;还要配置重试次数、最长重试时间(最大事件年龄),以及最终投递失败时保存事件的死信队列。

# 消息为什么会重复

消费者已经写入数据库,但删除消息前网络断了;可见性超时后,其他消费者再次收到这条消息。因此应该按至少一次的处理模型设计。

Visibility Timeout(可见性超时)表示消息被领取后暂时对其他消费者隐藏,不是消息 TTL,也不是分布式锁的永久所有权。耗时超过窗口仍可能出现并发处理;处理完才删除,长任务需要合适的续期或执行方案。

FIFO 支持分组内顺序和发送去重,但不能把 “外部写入 + 删除消息” 自动变成一个事务。使用 FIFO 仍要保护数据库、支付和外部 API 的业务幂等。

# 幂等键放在哪里

事件携带稳定 eventId、schemaVersion、业务 ID 和发生时间。相同业务结果重发应保持同一标识,不能每次重试生成新 UUID 让去重失效。

只涉及一个数据库时,可把 “插入已处理事件记录” 和 “业务状态更新” 放在同一事务:唯一键冲突表示已经处理。涉及外部付款或链上广播时,还要使用外部系统幂等能力、Outbox、状态机和对账,不能只靠内存 Set。

# 为什么先查再写还不够

消费者 A 和 B 同时收到重复事件,都查询到 “还没处理” ,随后都记账,就会发生重复效果。真正的互斥点应是数据库唯一约束与事务,不是应用里的先查后写。

下面是可在 PostgreSQL 独立连接执行的教学 SQL,仅建立会话级临时表。它展示同一个事件重复投递两次,最终只累计一次;示例不连接 AWS。

-- 临时表只在当前连接存在;生产表还需租户、版本及保留策略。
CREATE TEMP TABLE note_processed_events (event_id text PRIMARY KEY);
CREATE TEMP TABLE note_totals (account_id text PRIMARY KEY, amount numeric NOT NULL);
INSERT INTO note_totals VALUES ('alice', 0);

BEGIN;
-- 只有成功取得事件唯一键的处理者,才有资格修改业务结果。
WITH accepted AS (
  INSERT INTO note_processed_events VALUES ('evt-1')
  ON CONFLICT DO NOTHING RETURNING event_id
)
UPDATE note_totals SET amount = amount + 10
WHERE account_id = 'alice' AND EXISTS (SELECT 1 FROM accepted);
COMMIT;

BEGIN;
-- 模拟重投:唯一键已存在,accepted 为空,不会第二次加 10。
WITH accepted AS (
  INSERT INTO note_processed_events VALUES ('evt-1')
  ON CONFLICT DO NOTHING RETURNING event_id
)
UPDATE note_totals SET amount = amount + 10
WHERE account_id = 'alice' AND EXISTS (SELECT 1 FROM accepted);
COMMIT;
SELECT amount FROM note_totals WHERE account_id = 'alice'; -- 预期为 10。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25

生产中还要验证目标业务行确实存在、更新行数符合预期;不满足时必须回滚,不能只留下已处理标记。相同事件 ID 对应不同载荷也应拒绝并告警,不能悄悄吞掉数据冲突。出现事务失败或序列化冲突时整体重试,事务提交后才确认消息。

# 四个故障位置决定恢复方式

故障发生时 持久化事实 恢复动作
数据库事务提交前 去重标记和业务结果都应回滚 重新消费
提交后、消息确认前 业务已完成,消息可能再出现 唯一键阻止重复业务效果,再确认
API 已提交任务、发消息前 若无 Outbox,任务可能永远不被执行 在同一事务保存待投递记录,由投递器补发
外部付款已发出、响应超时 外部结果未知,本地事务无法撤销付款 使用对方幂等键查询或重试,必要时对账

Outbox 投递器本身也可能在发出消息后、标记已发前崩溃,因此 Outbox 解决不漏投递,不保证绝不重投;消费者幂等仍然必要。概念解释与项目落地复用 最终一致性与对账。

# 重试不能把故障越放越大

可见性窗口应覆盖处理时间并留出余量,长任务必要时续期,但续期失败仍可能导致并发执行,不能替代幂等。数据库故障时增加退避并限制消费者并发;数据格式永远不合法时尽快隔离,而不是反复占用工作线程。

重放 DLQ 前先修好根因,再选少量消息验证业务结果和重复保护,逐步放量。观察消息最老年龄、成功处理速率和下游错误,不能只看队列长度下降——消息被丢掉也会让长度下降。

# 部分批次失败怎样配合

一批十条消息中一条失败,默认整批重试会重复处理其他九条。Lambda 消费 SQS 时,需要同时配置 ReportBatchItemFailures 并返回失败的 messageId;只改代码或只改配置都不完整。

下面是可直接运行的本地 .mjs 教学示例,演示标准队列的批次协议,不调用 AWS。Map 只用于展示幂等判断,不能作为生产持久化去重。

import assert from 'node:assert/strict';

const results = new Map(); // 本地样本结果;生产应换成事务化持久存储。

async function processRecord(record) {
  const event = JSON.parse(record.body);
  if (typeof event.eventId !== 'string' || !event.eventId ||
      typeof event.report !== 'string') {
    throw new Error('INVALID_EVENT'); // 本例选择失败并等待后续隔离处理。
  }
  if (results.has(event.eventId)) return; // 同一业务事件重发,不再产生第二份结果。
  results.set(event.eventId, event.report);
}

export async function handler(event) {
  const batchItemFailures = [];
  for (const record of event.Records) {
    try {
      await processRecord(record);
    } catch {
      // itemIdentifier 必须是 SQS 的 messageId,而不是自定义 eventId。
      batchItemFailures.push({ itemIdentifier: record.messageId });
    }
  }
  return { batchItemFailures };
}

const valid = JSON.stringify({ eventId: 'report-7-v1', report: '完成' });
const response = await handler({ Records: [
  { messageId: 'm1', body: valid },
  { messageId: 'm2', body: valid }, // 重复业务事件,不重复写入。
  { messageId: 'm3', body: 'invalid-json' },
] });
assert.equal(results.size, 1);
assert.deepEqual(response, { batchItemFailures: [{ itemIdentifier: 'm3' }] });
console.log('批次与去重示例通过');
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36

生产中要区分永久错误与暂时失败。协议非法的消息应记录原因并隔离,策略允许时确认消费;数据库暂时不可用则重试。FIFO 要保持顺序,首次失败后还需正确报告相关未处理消息,不能机械照搬标准队列 “继续处理整批” 的方式。

# DLQ 不是自动修复器

DLQ(Dead-Letter Queue,死信队列)保存超过重试策略仍失败的消息。先查错误、修复原因、确认消费者幂等,再小批量 redrive(重新送回处理)。不排查就全部重放,可能重复扣款或再次把系统压垮。

人工从队列 receive 消息会影响可见性和接收计数,不能当作完全无副作用的查看。日常先看指标和脱敏日志,取样需要明确操作范围。

原项目还有一次 “Topic 接受消息,但队列收不到” 的故障:SQS 的 AWS 托管 KMS key 无法补充所需服务权限。最终采用 SSE-SQS;若确需自管 KMS,则分别检查队列策略与密钥策略,并评估费用。加密配置不是只打开一个开关。

# 面试时可以这样回答

用了消息队列,怎样避免重复执行?参考答案

我按至少一次处理设计,让同一个业务事件有稳定幂等键。在数据库事务里同时检查去重和更新业务,成功后才确认消息;外部付款等副作用还要结合对方的幂等能力和对账。批量消费只重试失败项,持续失败进入 DLQ,修复后再有控制地重放。

数据库提交和消息确认不能放在一起,怎样保证正确?参考答案

我不假设两者能原子完成,而是先在数据库事务里一起写去重记录和业务结果,再确认消息。提交后确认失败会导致重投,但唯一键能阻止第二次业务效果。发送侧用 Outbox 保证任务和待发送事件同事务保存;外部付款再结合对方幂等键和对账,这样每个失败窗口都有恢复办法。

# 复用来源与官方参考

整理自原项目稳定性与异步事件章节,以及性能 Processor 的重试经验。