Kafka消息丢了怎么办?从重试到DLQ的完整可靠性设计
SryYa Three

可靠性问题

我们的业务在引入消息队列后,它确实带了给我们便利,但我们也就不得不面对它带来的问题:发送消息失败怎么处理呢?特别是我们的业务跟金融相关,消息处理不当就会给我们和客户带来巨大损失。因此解决消息队列的可靠性问题迫在眉睫。

今天,我们就以电商秒杀 + 消息队列Kafka为场景讲解如何处理消息发送失败以及死信队列的使用。


发送失败处理

什么时候需要失败处理

在项目中,并不是什么地方都需要失败处理(比如一些可有可无的日志),什么时候需要发送失败处理呢?

我们可以问自己几个问题:

  1. 这条消息丢了影响大吗?

    大 → 必须处理

    无所谓 → 可以放宽

  2. 有没有补偿手段?

    有(例如数据库消息表)→ 有其他方案可以兜底,失败处理优先级可以降低。

    没有 → 必须保证发送成功

  3. 是否在事务之后?

    是 → 高危场景,必须防止数据不一致问题。

如果消息队列涉及到业务流程中的关键节点,失败处理势在必行。

失败处理设计

接下来,我们就以一段处理缓存失效通知失败的代码为例,讲解发送失败处理的思路。

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
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
/**
* 发送失败处理:结构化日志、上报错误指标、自适应退避重试、超限后转交 DLQ。
*/
@Override
protected void afterSendFailure(
final String topic, //发送失败的topic
final MessageExtend<SeckillVoucherInvalidationMessage> message, //发送失败的消息
final Throwable throwable //异常信息
) {
// 结构化日志,便于定位问题
final SeckillVoucherInvalidationMessage body = message.getMessageBody();
final Long voucherId = body.getVoucherId();
final String reason = body.getReason();
final String errMsg = throwable == null ? "unknown" : throwable.getMessage();
log.error("SeckillVoucherInvalidation send failed, topic={}, uuid={}, key={}, voucherId={}, reason={}, error= {}",
topic, message.getUuid(), message.getKey(), voucherId, reason, errMsg, throwable);

//如果是死信队列发过来的消息,直接上报错误指标,不进行重试
if (topic.contains(DLQ)) {
safeInc("seckill_invalidation_dlq", "reason", "send_failures");
return;
}else {
// 指标:失败计数
safeInc("seckill_invalidation_send_failures", "topic", topic);
}

//失败重试,通过 header 标记重试次数,避免无限重试
Map<String, String> headers = message.getHeaders();
headers = headers == null ? new HashMap<>(8) : new HashMap<>(headers);
int retryCount = 0;
try {
// 取出重试次数RETRY_COUNT
if (headers.containsKey(RETRY_COUNT)) {
retryCount = Integer.parseInt(headers.get(RETRY_COUNT));
}
} catch (Exception ignore) {
}

// 重试次数 < 最大重试次数限制 --> 退避重试(随着重试次数提升间隔时间)
if (retryCount < retryMaxAttempts) {
long backoff = Math.min(initialBackoffMillis * (1L << retryCount), maxBackoffMillis);
// 更新重试次数
headers.put(RETRY_COUNT, String.valueOf(retryCount + 1));
headers.put("lastError", truncate(errMsg));
message.setHeaders(headers);
log.warn("Retry sending cache invalidation, topic={}, uuid={}, voucherId={}, retryCount={}, backoffMs={}",
topic, message.getUuid(), voucherId, retryCount + 1, backoff);
// 上报监控中心,讲明"什么发生了重试"
safeInc("seckill_invalidation_send_retries", "topic", topic);
// 退避,等待重试
sleepQuietly(backoff);
// 异步发送消息重试:若再失败会再次进入本方法,直到超过最大重试次数
sendRecord(topic, message);
return;
}

// 重试次数 > 最大重试次数限制 --> 进入死信队列
// 放入 DLQ,便于后续人工/自动补偿
final String dlqReason = "send_invalid_cache_broadcast_failed: " + truncate(errMsg);
try {
// 发送到死信队列
sendToDlq(topic, body, dlqReason);
log.warn("Send cache invalidation to DLQ, originalTopic={}, uuid={}, voucherId={}, dlqReason={}",
topic, message.getUuid(), voucherId, dlqReason);
auditLog.warn("DLQ_PUBLISH|topic={}|uuid={}|key={}|voucherId={}|reason={}",
topic, message.getUuid(), message.getKey(), voucherId, dlqReason);
// 上报监控中心,记录消息发送到死信队列
safeInc("seckill_invalidation_send_dlq", "topic", topic);
} catch (Exception e) {
log.error("Send cache invalidation to DLQ failed, originalTopic={}, uuid={}, voucherId={}, error={}",
topic, message.getUuid(), voucherId, e.getMessage(), e);
safeInc("seckill_invalidation_send_dlq_failures", "topic", topic);
}
}

这一段代码的流程为:

image


重试机制

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
if (retryCount < retryMaxAttempts) {
long backoff = Math.min(initialBackoffMillis * (1L << retryCount), maxBackoffMillis);
// 更新重试次数
headers.put(RETRY_COUNT, String.valueOf(retryCount + 1));
headers.put("lastError", truncate(errMsg));
message.setHeaders(headers);
log.warn("Retry sending cache invalidation, topic={}, uuid={}, voucherId={}, retryCount={}, backoffMs={}",
topic, message.getUuid(), voucherId, retryCount + 1, backoff);
// 上报监控中心,讲明"什么发生了重试"
safeInc("seckill_invalidation_send_retries", "topic", topic);
// 退避,等待重试
sleepQuietly(backoff);
// 异步发送消息重试:若再失败会再次进入本方法,直到超过最大重试次数
sendRecord(topic, message);
return;
}

在这里,我们使用了一个简单的退避重试:通过Header记录的重试次数,动态地计算等待时间,避免重试带来的压力。

兜底方案

如果遇到消息队列彻底挂了这种极端情况,连重试和死信队列都发不出去怎么办?

先别慌,还有兜底办法。那就是本地消息表:

先将消息持久化到 MySQL,由定时任务补偿发送。

下面给出一种在秒杀时的实现思路:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
@Transactional
public void handleSeckill(Voucher voucher) {
// 执行秒杀逻辑
voucherMapper.updateStock(voucher);

// 构造消息记录并保存到本地表
OutboxMessage message = new OutboxMessage(voucher, "seckill_topic");
outboxMapper.insert(message);

// 事务提交后,手动异步发送
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
asyncTask.sendToKafka(message);
}
});
}

流程:

原子写入: 由于我们开启了事务,所以执行秒杀逻辑和插入本地消息表保证了其原子性。

异步调用: 接下来,事务提交后,立即异步调用 Kafka 发送。如果发送成功,更新消息表状态为“成功”或直接删除。

定时兜底: 我们还可以再加一重保险:启动一个定时任务,每分钟扫描一次 msg_outbox 中状态为“待发送”且超时的记录,重新发送。

消息幂等: 需要注意的是,由于定时任务可能会重复发送,消费者端必须根据消息 UUID 实现幂等逻辑。

消息表的设计: 对于表msg_outbox,字段包括:id, payload,topic, status (待发送/成功), retry_count, create_time,可以根据业务需求自行修改。


死信队列

对于进入死信队列的消息,其实已经是进行过重试的消息了,所以再进行重试已经意义不大了。

我们能做的就是人工补偿、记录错误信息和上报指标了。

死信队列的Topic可以设计成”原业务Topic_DLQ”的样式,更加清晰易懂。

那么死信队列的思路就很简单了,下面见代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Override
protected void doConsume(MessageExtend<SeckillVoucherInvalidationMessage> message) {
// 获取消息
SeckillVoucherInvalidationMessage body = message.getMessageBody();
// 记录日志并上报监控中心
if (Objects.isNull(body.getVoucherId())) {
log.warn("DLQ消息载荷为空或voucherId缺失, uuid={}", message.getUuid());
safeInc("seckill_invalidation_dlq_replay_skipped", "reason", "invalid_payload");
return;
}

safeInc("seckill_invalidation_dlq", "reason", "invalid_payload");

auditLog.error("SECKILL_INVALIDATION_DLQ | message={}", JSON.toJSONString(message));
}

总结

死信队列不是终点,而是问题处理的起点。

引入消息队列虽然解耦了业务,但也带来了可靠性的挑战。通过本文的实践,我们可以总结出保障消息可靠性的三道防线

  1. 第一道防线:自适应重试。 利用指数退避算法,在网络抖动时给系统留出喘息机会,避免盲目重试加剧负载。
  2. 第二道防线:死信队列。 当重试耗尽时,将问题消息隔离并记录审计日志,能有效防止核心业务数据的静默丢失。
  3. 第三道防线:可观测性。 没有监控的系统是盲目的。通过结构化日志和多维度的指标上报,我们要做到常说的**“感知早于投诉”**。

虽然重试和 DLQ 解决了大部分问题,但在极端的金融级场景下,我们还可以引入 本地消息表 来实现数据库与 MQ 的强一致性。技术方案没有银弹,根据业务价值,选择最合适的可靠性等级,才是后端开发的进阶之道。

由 Hexo 驱动 & 主题 Keep
总字数 28.1k 访客数 访问量