Reactor Kafka并发消费时如何实现消息失败重试不丢失
问题背景
- 现有Kafka消息流转项目核心逻辑:消费入站Topic消息,完成业务转换处理后发送至出站Topic,要求基础设施故障导致出站Topic不可用、消息发送报错时,仅提交最后一条处理并发送成功的消息偏移量,从发送失败位置重新开始消费,严格保证消息不丢失。
- 非响应式Kafka实现中,该逻辑可通过
ack.nack(index, sleep)方法轻松实现,已有对应@KafkaListener注解式批量消费实现代码。 - 改造目标:将逻辑迁移为响应式实现,获得更优性能与原生背压支持,同时实现不同分区的消息并行处理。
当前实现异常现象
- 初版实现参考Reactor Kafka官方文档编写:使用
groupBy按TopicPartition分组、BoundedElastic类型Scheduler实现分区并行处理,结合KafkaSender发送消息、Retry.backoff配置退避重试策略。 - 测试场景:将出站Topic修改为无权限的不存在Topic模拟发送故障
- 异常表现:
- 消息发送失败后代码未按预期重试,新消息到来时会直接丢弃之前发送失败的消息,造成消息丢失
- 测试日志中出现调度器工作线程未捕获的空指针异常
核心诉求
- 保证消息零丢失,实现发送失败时从失败偏移量重试消费的逻辑
- 提供相关性能优化建议
- 所有代码片段、API名称、日志内容、技术术语均保持原有形式
根因说明
初版实现出现消息丢失、空指针异常的核心原因有3个:
- 偏移量提交逻辑失控:Reactor Kafka默认自动提交逻辑是按流的请求进度推进,若重试逻辑配置在流的最外层、或错误信号没有被正确隔离在单条消息处理链路内,一旦发送报错触发流重启,会从最后一次已提交的偏移量开始拉取,未发送成功且未提交的消息会被直接丢弃。
- 分区处理顺序被打乱:使用
BoundedElastic做并行时如果没有做线程亲和性绑定,同分区的消息可能被调度到不同线程处理,既破坏了分区内消息的顺序性,还会因为流取消时偏移量对象被提前释放,触发工作线程的空指针异常。 - 重试作用域错误:如果把
Retry.backoff配置在整个消费流的外层,单条消息发送失败会触发整个流的重新订阅,此时所有分区的消费进度都会被重置,不仅会重复消费其他分区已经成功的消息,还可能因为异步提交的偏移量已经落盘导致部分失败消息被跳过。
修复实现方案
实现必须严格遵守4个约束:单分区内严格串行处理、重试逻辑绑定在单条/单批次消息发送链路内、手动控制偏移量提交时机、发送失败时阻塞当前分区消费直到成功且不影响其他分区并行。
核心实现代码如下:
// 消费端配置:关闭自动偏移量提交 ReceiverOptions<String, String> receiverOptions = ReceiverOptions.<String, String>create(consumerProps) .subscription(Collections.singletonList(inboundTopic)) .autoCommit(false) .commitInterval(Duration.ZERO) .commitBatchSize(0); KafkaReceiver.create(receiverOptions) .receive() // 按TopicPartition分组,同分区消息路由到同一个独立Flux .groupBy(record -> new TopicPartition(record.topic(), record.partition())) .flatMap(partitionFlux -> partitionFlux // 同分区消息绑定专属线程,保证线程亲和性,全程不做跨线程切换 .publishOn(Schedulers.newBoundedElastic(1, 1024, "kafka-pt-" + partitionFlux.key().partition()), 1) // 同分区按偏移量顺序串行处理,不打乱消费顺序 .concatMap(record -> { // 执行业务转换逻辑 String transformedValue = businessConvert(record.value()); SenderRecord<String, String, Object> sendRecord = SenderRecord.create( outboundTopic, null, null, record.key(), transformedValue, record.receiverOffset() ); return kafkaSender.send(Mono.just(sendRecord)) // 重试逻辑绑定在单条消息发送链路内,失败不终止分区流 .retryWhen(Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(1)) .maxBackoff(Duration.ofSeconds(30)) .filter(ex -> ex instanceof KafkaException || ex.getCause() instanceof RetriableException) ) // 仅当消息发送成功后,才提交当前消息的偏移量 .then(Mono.fromRunnable(() -> record.receiverOffset().acknowledge())) // 异常场景下持续重试当前消息,不拉动同分区下一条消息 .onErrorResume(ex -> { log.error("Partition {} offset {} send failed, retrying", partitionFlux.key(), record.offset(), ex); return Mono.empty(); }); }, 1) // flatMap并发度配置为和分区数一致,保证不同分区完全并行 ).subscribe();
该实现下,单条消息发送失败时会持续重试当前消息,不会拉取同分区后续消息,也不会提前提交任何偏移量,从根本上避免消息丢失;同时因为每个分区绑定独立线程,不会出现偏移量对象被提前释放的空指针问题。
性能优化建议
- 单分区内可通过
windowTimeout(200, Duration.ofMillis(50))做小批量聚合,攒够200条消息或50ms时间窗口就批量发送,大幅减少网络IO次数,批量发送成功后提交批次内最大的偏移量即可。 - 调度器线程数配置为和入站Topic总分区数一致即可,不要创建过多线程导致不必要的上下文切换开销。
- 重试策略中可加入连续失败阈值告警,避免长时间基础设施故障时线程空转。
- 发送端配合配置
linger.ms=50、batch.size=16384参数,和消费端攒批逻辑适配,进一步提升发送吞吐量。 - 同分区处理链路内禁止使用无顺序保证的
flatMap、跨线程异步传递等操作,避免偏移量乱序提交。
内容的提问来源于stack exchange,提问作者zeeshan
相关产品推荐
相关产品推荐

