You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.31 23:57:35