本节把 TaskHub 推进到「重复也无害」:8.1 的至少一次投递意味着同一条消息可能被处理两遍,本节给消费者加上幂等去重、指数退避重试与死信队列,让坏消息既不丢也不无限重试。
适用版本:Go 1.27(实测go1.27.0),Redis 7、PostgreSQL 17(容器)。
8.2 幂等、重试与死信
8.1 节选了「至少一次」投递,代价写在明面上:消费者崩溃后消息会被重投,同一条消息可能被处理多次。异步链路要能用,就必须让重复处理不产生副作用——这就是幂等。本节把幂等、重试、死信这三件配套的事一次做齐。
8.2.1 幂等的本质:把「执行」变成「有条件执行」
幂等的意思是:同一个操作执行一次和执行多次,结果相同。对消费场景,最常见的实现是给每条消息一个唯一 ID,处理前先检查「这个 ID 是不是处理过」,处理过就跳过。
TaskHub 的事件里天然带一个 event_id(生产者生成,通常用 UUID)。去重有两种落点:
| 落点 | 手段 | 特点 |
|---|---|---|
| 缓存 | Redis SETNX processed:{id} | 快,但 Redis 可能丢数据,重启后去重记录没了 |
| 数据库 | 唯一约束 + ON CONFLICT DO NOTHING | 持久、可靠,是真正的真相来源 |
实践里两道都上:Redis 挡掉绝大多数重复(快),数据库唯一约束兜底(可靠)。只靠 Redis 的话,Redis 一旦被清空或故障切换,历史重复就挡不住了。
8.2.2 第一道:Redis 快速去重
用 SETNX(set if not exists)尝试写入去重键,写入成功说明是首次,失败说明已处理过:
dedupKey := "processed:" + eventID
first, _ := rdb.SetNX(ctx, dedupKey, "1", 24*time.Hour).Result()
if !first {
// 已处理过,直接 ACK 跳过
rdb.XAck(ctx, "tasks:events", "workers", msgID)
return
}
TTL 设成「远大于消息可能被重投的窗口」(比如 24 小时)即可,不需要永久——超过这个窗口的消息早就被死信或人工处理了。
8.2.3 第二道:数据库唯一约束
Redis 之外,用数据库的唯一约束做最终仲裁。建一张「已处理事件」表,event_id 做主键,插入时用 ON CONFLICT DO NOTHING:插入影响 0 行说明是重复,影响 1 行说明是首次。
CREATE TABLE processed_events (
event_id text PRIMARY KEY,
handled_at timestamptz DEFAULT now()
);
func handle(ctx context.Context, pool *pgxpool.Pool, eventID string) (bool, error) {
tag, err := pool.Exec(ctx,
`INSERT INTO processed_events(event_id) VALUES($1)
ON CONFLICT (event_id) DO NOTHING`, eventID)
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil // true=首次处理
}
同一个 event_id 处理三次,实测只有第一次真正执行:
$ go run ./ch8/idem
第 1 次: 首次处理,执行业务副作用
第 2 次: 已处理过,跳过(幂等)
第 3 次: 已处理过,跳过(幂等)
processed_events 行数 = 1
关键点:插入去重记录与业务副作用要在同一个事务里——否则「插入了去重记录但副作用没做完」或「副作用做了但去重记录没落库」都会出问题。用事务包起来,要么全成、要么全不成,重投时能安全重来。
把去重记录与副作用放进同一个事务,实测三种情况:
tx, _ := pool.Begin(ctx)
defer tx.Rollback(ctx)
tag, _ := tx.Exec(ctx,
`INSERT INTO processed_events(event_id) VALUES($1) ON CONFLICT DO NOTHING`, eventID)
if tag.RowsAffected() == 0 {
return nil // 重复,幂等跳过
}
if _, err := tx.Exec(ctx, `INSERT INTO notifications(event_id) VALUES($1)`, eventID); err != nil {
return err // 副作用失败,去重记录一并回滚
}
return tx.Commit(ctx)
$ go run ./ch8/txidem
副作用失败: err=副作用失败,整体回滚
回滚后 processed_events=0 notifications=0
重投成功: err=<nil>
提交后 processed_events=1 notifications=1
重复投递后 notifications=1(未重复)
副作用失败时事务整体回滚,processed_events 里没有留下去重记录——所以重投时能重新处理,不会被自己的去重记录挡住。这正是「去重与副作用同事务」的价值:如果去重记录先单独提交了,副作用失败后的重投就会被误判成重复而跳过,导致业务永远做不成。
8.2.4 重试:指数退避
处理失败时不能立刻无限重试——如果是下游临时抖动,退避一会儿就好了;如果是永久错误,快速重试只是浪费。标准做法是指数退避 + 上限:
func backoff(attempt int) time.Duration {
base := 20 * time.Millisecond
d := base << attempt // 20, 40, 80, 160...
if d > time.Second {
d = time.Second // 封顶
}
return d
}
第一次失败退避 20ms、第二次 40ms、第三次 80ms,指数增长。实测一条注定失败的消息重试三次的过程:
$ go run ./ch8/dlq
msg-2 第 1 次失败(downstream 500),退避 20ms
msg-2 第 2 次失败(downstream 500),退避 40ms
msg-2 第 3 次失败(downstream 500),退避 80ms
msg-2 超过重试上限 -> 进死信
生产里 base 通常从几百毫秒起、封顶到几十秒,并加随机抖动(jitter),避免大量消息在同一时刻一起重试(惊群)。退避还有个变体是服务端返回 Retry-After——如果下游明确告诉你「1 分钟后再来」,就听它的,别自己瞎退避。
8.2.5 什么能重试,什么不能
盲目重试会放大故障。先给错误分类:
| 错误类型 | 例子 | 能否重试 |
|---|---|---|
| 瞬时/网络 | 连接超时、5xx、限流 | 能,退避后重试 |
| 依赖暂时不可用 | 下游维护中 | 能,退避后重试 |
| 业务校验失败 | 参数非法、状态不允许 | 不能,重试多少次都一样,直接进死信 |
| 数据不存在 | 引用的实体被删 | 视语义,通常不能 |
| 幂等冲突 | 已处理过 | 不算错误,跳过 |
判据是:重试能不能改变结果。不能改变的,重试是纯粹的浪费,应该立刻判死信。给错误打上「可重试」标记(如自定义 RetryableError 类型)是常见的工程做法,让重试逻辑一眼可判。
8.2.6 死信队列:坏消息的收容所
重试超过上限的消息不能再无限循环,要把它们移出主链路,送进死信队列(DLQ, Dead Letter Queue):
if err != nil {
rdb.XAdd(ctx, &redis.XAddArgs{
Stream: "tasks:events:dlq",
Values: map[string]any{"id": eventID, "reason": err.Error()},
})
}
rdb.XAck(ctx, "tasks:events", "workers", msgID) // 从主链路移除
实测一条永远失败的消息进了死信:
$ go run ./ch8/dlq
死信条数 = 1
DLQ 1791598968010-0 map[id:msg-2 reason:downstream 500]
死信的设计要点:
- 保留原因:
reason字段记录最后一次失败的错误,是排查的唯一线索。 - 保留原始消息:死信里要能还原出完整的事件内容(这里简化为
id,生产里应带完整 payload)。 - 从主链路 ACK 掉:否则这条消息还会被
XAUTOCLAIM重投,形成「重试 → 死信 → 又重投」的循环。
8.2.7 死信的后续:不是终点
死信队列不是垃圾桶,它需要有人管:
- 告警:DLQ 一旦非空就告警(第 10 章),因为每一条死信都代表一次业务失败。
- 可视化:把死信内容做成可查的页面,让运维能看「哪条消息、为什么失败」。
- 重放:修完根因后,把死信重新投回主 stream 重放。重放前要确认修复真的生效,否则只是再进一次死信。
- 定期清理:长期无人处理的死信要归档或删除,别让 DLQ 也无限增长。
TaskHub 的约定是:DLQ 非空触发告警,on-call 在值班手册里查处理流程(第 16 章)。没有后续处理的死信队列等于没有死信队列——消息只是从「重试循环」挪到了「无人问津」。
8.2.8 幂等做在业务层还是框架层
去重逻辑放哪儿,是个架构选择:
| 位置 | 做法 | 优点 | 缺点 |
|---|---|---|---|
| 业务层 | 每个 handler 自己判重 | 精确、能表达业务语义 | 重复代码多、易漏 |
| 框架层 | 消费框架统一拦截 event_id | 一致、不漏 | 对「什么算同一事件」的理解固化 |
TaskHub 的折中是:框架层提供 WithIdempotency(eventID, fn) 包装器,业务层决定去重键与 TTL:
func WithIdempotency(ctx context.Context, eventID string, fn func() error) error {
first, _ := rdb.SetNX(ctx, "processed:"+eventID, "1", 24*time.Hour).Result()
if !first {
return nil // 已处理
}
if err := fn(); err != nil {
rdb.Del(ctx, "processed:"+eventID) // 失败回滚去重标记,允许重试
return err
}
return nil
}
注意包装器在 fn 失败时删掉 Redis 去重标记——否则失败的消息因为标记还在,重投时会被直接跳过,永远不重试。这个「失败即回滚去重标记」的细节和数据库事务回滚是同一个道理,两种落点都要遵守。
8.2.9 常见坑
- 只做 Redis 去重:Redis 故障或清空后重复全部漏过,数据库唯一约束才是真相。
- 去重记录与业务副作用不在同一事务:半成功状态在重投时会出问题。
- 无限重试不封顶:坏消息永远占着消费者,把队列拖死。
- 退避不加抖动:大量消息同时重试,把刚恢复的下游再次打垮。
- 重试业务校验失败:参数错了重试一万次还是错,应直接判死信。
- 死信不进不 ACK:消息既在 DLQ 又在主链路 pending,被反复重投。
- 幂等键用「内容哈希」:同一事件内容变了就当成新消息,去重失效,应该用生产者生成的稳定
event_id。 - 死信无人处理:没有告警和重放机制,坏消息静默堆积。
小结
- 至少一次投递必然带来重复,消费端幂等是配套要求,不是可选项。
- 去重两道防线:Redis
SETNX快速挡(快),数据库唯一约束兜底(可靠)。 - 幂等落库用
INSERT ... ON CONFLICT DO NOTHING,靠RowsAffected判断是否首次,实测三次处理只执行一次。 - 重试用指数退避 + 上限 + 抖动;能否重试取决于「重试能不能改变结果」。
- 超过上限的消息进死信队列,保留原因、从主链路 ACK,并配套告警与重放。
到这里 TaskHub 的消息链路能可靠投递、幂等消费、失败重试、坏消息收容了。但还有一类任务既不是「收到事件就做」,也不是「立刻做」——它们要定时触发或延迟到某个时刻才执行。下一节讲定时任务与延迟队列。
阅读导航:上一节:8.1 生产者/消费者与可靠投递 · 下一节:8.3 定时任务与延迟队列 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。