Spring Kafka:如何通过kafkaTemplate实现同分区同偏移量失败重试
基于KafkaTemplate实现发送失败同分区同偏移量重试方案
首先明确核心逻辑:Kafka中只有被Broker成功确认写入日志的消息才会分配正式的分区和偏移量,发送阶段就失败的消息根本不会被Broker记录,不存在“预留偏移量”的说法。你要实现的「相同分区、相同偏移量位置重试」,本质是保证重试消息和原消息路由到完全一致的分区,且不会被其他消息插队,让重试消息成为对应分区下一个被追加写入的消息,自然就会落到原消息预期的偏移量位置。
具体实现步骤
1. 前置配置调整
首先关掉Kafka生产者自带的默认异步重试,避免内部重试打乱消息顺序,导致偏移量错位:
- 在配置文件中将
spring.kafka.producer.retries设为0,禁用原生自动重试 - 保证生产端使用的分区器逻辑固定:如果是按Key哈希路由分区,要保证相同Key永远计算出相同分区号,不要使用随机分区策略;如果是手动指定分区发送,重试时直接复用原分区号即可
- 给KafkaTemplate配置发送失败监听器,捕获发送失败的消息上下文
配置示例代码:
@Configuration public class KafkaProducerConfig { @Bean public KafkaTemplate<String, byte[]> kafkaTemplate(ProducerFactory<String, byte[]> producerFactory) { KafkaTemplate<String, byte[]> template = new KafkaTemplate<>(producerFactory); template.setProducerListener(new ProducerListener<String, byte[]>() { @Override public void onError(ProducerRecord<String, byte[]> record, RecordMetadata metadata, Exception exception) { // 捕获失败消息,投递到重试队列 retryManager.addFailedRecord(record, metadata); } }); return template; } }
2. 实现分区维度的串行重试逻辑
要保证偏移量位置正确,核心是同分区的重试消息必须按原发送顺序串行发送,等上一条消息收到Broker的成功ACK后再发下一条,避免多线程并发发送导致消息插队。
核心实现逻辑:
- 为目标Topic的每个分区单独维护一个阻塞重试队列,每个队列对应一个独立的工作线程,保证同分区消息串行发送
- 重试时构造和原消息完全一致的
ProducerRecord,复用原消息的Topic、分区号、Key、Value、消息头、时间戳,保证路由结果和原消息完全一致 - 采用同步发送模式,发送成功再处理队列下一条消息,发送失败可按业务配置重试次数,超过阈值转死信处理
重试逻辑示例代码:
@Component public class KafkaRetryManager { @Autowired private KafkaTemplate<String, byte[]> kafkaTemplate; // 按分区号维护独立的重试队列,保证同分区消息顺序 private final Map<Integer, BlockingQueue<ProducerRecord<String, byte[]>>> retryQueues = new ConcurrentHashMap<>(); private static final int MAX_RETRY_COUNT = 3; @PostConstruct public void startRetryWorkers() { // 初始化时获取目标Topic的分区数,为每个分区启动独立重试线程 List<PartitionInfo> partitions = kafkaTemplate.partitionsFor("你的业务Topic名称"); for (PartitionInfo partition : partitions) { int partitionId = partition.partition(); retryQueues.put(partitionId, new LinkedBlockingQueue<>()); new Thread(() -> doRetry(partitionId), "kafka-retry-thread-" + partitionId).start(); } } public void addFailedRecord(ProducerRecord<String, byte[]> failedRecord, RecordMetadata metadata) { Integer targetPartition = metadata != null ? metadata.partition() : failedRecord.partition(); // 如果发送失败时没拿到分区信息,就用和生产端一致的分区逻辑计算目标分区 if (targetPartition == null) { targetPartition = calculatePartition(failedRecord); } // 构造和原消息完全一致的重试记录 ProducerRecord<String, byte[]> retryRecord = new ProducerRecord<>( failedRecord.topic(), targetPartition, failedRecord.timestamp(), failedRecord.key(), failedRecord.value(), failedRecord.headers() ); retryQueues.get(targetPartition).add(retryRecord); } private void doRetry(int partitionId) { BlockingQueue<ProducerRecord<String, byte[]>> queue = retryQueues.get(partitionId); while (true) { try { ProducerRecord<String, byte[]> record = queue.take(); int retryCount = 0; boolean sendSuccess = false; while (retryCount < MAX_RETRY_COUNT && !sendSuccess) { try { // 同步发送,等待Broker确认 kafkaTemplate.send(record).get(3, TimeUnit.SECONDS); sendSuccess = true; } catch (Exception e) { retryCount++; // 可以加间隔退避逻辑,避免频繁重试打垮Broker Thread.sleep(100L * retryCount); } } if (!sendSuccess) { // 超过最大重试次数,转死信队列或者落库告警,避免消息丢失 handleDeadLetterMsg(record); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } // 必须和生产端使用的分区器计算逻辑完全一致,保证同Key路由到相同分区 private int calculatePartition(ProducerRecord<String, byte[]> record) { // 这里替换成你实际使用的分区计算逻辑 return Math.abs(record.key().hashCode()) % retryQueues.size(); } private void handleDeadLetterMsg(ProducerRecord<String, byte[]> record) { // 自定义死信处理逻辑,比如落库、发告警 } }
关键注意事项
- 绝对不要依赖Kafka生产者自带的重试机制实现需求,原生重试是异步执行的,默认开启
enable.idempotence时虽然能保证不重复,但无法保证消息顺序,会出现后来的消息先发送成功的情况,导致偏移量错位。 - 内存队列存失败消息有丢数风险,如果业务对可靠性要求高,可以把失败消息先持久化到本地磁盘或者数据库,应用启动时先加载未完成的重试任务,避免重启丢消息。
- 重试时不要修改消息的Key、Value、Headers任何内容,否则可能出现路由分区不一致、业务逻辑异常的问题。
内容的提问来源于stack exchange,提问作者Rajeev Akotkar
相关产品推荐
相关产品推荐

