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

Kafka生产者retries=5配置下RecordTooLargeException重试及失败处理方案

问题根因

Kafka生产者内置的retries配置仅对可重试异常自动触发重试,org.apache.kafka.common.errors.RecordTooLargeException属于默认不可重试异常:该异常触发时,序列化后的消息大小已经超过客户端max.request.size阈值,请求根本不会发送到Broker,重复发送相同消息结果不会有任何变化,因此生产者会直接抛出异常交付回调,不会执行内置重试逻辑。

如果消息大小固定超过阈值且无调整逻辑,重试本身没有实际意义,需要先同步调整以下配置保证大小匹配:

  • 生产者端max.request.size
  • Broker全局配置message.max.bytes
  • Topic级别配置max.message.bytes

如果确实需要对该类异常做自定义重试(比如重试前压缩/拆分消息、动态调整配置),可以通过以下两种方案实现。

实现方案

方案1:基于Spring Retry实现声明式重试(适配Spring Kafka/KafkaTemplate场景)

该方案无侵入业务发送逻辑,通过AOP拦截异常自动触发重试,适合使用KafkaTemplate.addCallback发送消息的Spring生态项目。

  1. 引入依赖:Spring Boot项目直接引入spring-retry和spring-boot-starter-aop即可
  2. 添加重试配置类开启重试功能:
@Configuration
@EnableRetry
public class KafkaRetryConfig {
}
  1. 封装发送逻辑,配置重试规则与失败回调:
@Service
public class KafkaProducerService {

    private final KafkaTemplate<String, byte[]> kafkaTemplate;
    // 配置最大重试次数为5次
    private static final int MAX_RETRY_TIMES = 5;

    public KafkaProducerService(KafkaTemplate<String, byte[]> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    @Retryable(
            // 指定触发重试的异常类型
            retryFor = RecordTooLargeException.class,
            // maxAttempts包含首次请求,5次重试需要设置为6
            maxAttempts = MAX_RETRY_TIMES + 1,
            // 重试退避策略:间隔1s发起下一次重试,避免空转打满资源
            backoff = @Backoff(delay = 1000)
    )
    public void sendMessage(String topic, byte[] payload) {
        ListenableFuture<SendResult<String, byte[]>> sendFuture = kafkaTemplate.send(topic, payload);
        sendFuture.addCallback(new ListenableFutureCallback<>() {
            @Override
            public void onFailure(Throwable ex) {
                // 必须主动抛出需要重试的异常,否则Spring Retry无法感知失败触发重试
                if (ex instanceof RecordTooLargeException) {
                    throw (RecordTooLargeException) ex;
                }
                // 非重试异常直接走永久失败处理
                handlePermanentFailure(payload, ex);
            }

            @Override
            public void onSuccess(SendResult<String, byte[]> result) {
                // 自定义发送成功逻辑
                System.out.printf("消息发送成功,topic:%s, partition:%d, offset:%d%n",
                        result.getRecordMetadata().topic(),
                        result.getRecordMetadata().partition(),
                        result.getRecordMetadata().offset());
            }
        });
    }

    // 重试次数耗尽后自动触发该方法,入参需要和@Retryable方法的异常、参数顺序保持一致
    @Recover
    public void handleRetryExhausted(RecordTooLargeException ex, String topic, byte[] payload) {
        handlePermanentFailure(payload, ex);
    }

    private void handlePermanentFailure(byte[] payload, Throwable ex) {
        // 自定义永久失败逻辑:发送到DLT死信队列、写入数据库存储等
        System.out.printf("消息永久失败,异常信息:%s,准备执行兜底存储%n", ex.getMessage());
    }
}

方案2:原生客户端手动实现重试(无Spring依赖场景)

如果直接使用原生KafkaProducer,可以在回调层手动维护重试计数,实现重试逻辑:

public class ReliableKafkaProducer {

    private final KafkaProducer<String, byte[]> producer;
    private static final int MAX_RETRY_TIMES = 5;
    private static final long RETRY_BACKOFF_MS = 1000;

    public ReliableKafkaProducer(KafkaProducer<String, byte[]> producer) {
        this.producer = producer;
    }

    public void sendWithRetry(ProducerRecord<String, byte[]> record) {
        // 用原子整数计数,保证多线程场景下计数准确
        AtomicInteger retryCounter = new AtomicInteger(0);
        doSend(record, retryCounter);
    }

    private void doSend(ProducerRecord<String, byte[]> record, AtomicInteger retryCounter) {
        producer.send(record, (metadata, exception) -> {
            if (exception == null) {
                // 发送成功逻辑
                System.out.printf("消息发送成功,topic:%s, partition:%d, offset:%d%n",
                        metadata.topic(), metadata.partition(), metadata.offset());
                return;
            }

            // 判断是否满足重试条件:异常类型匹配+重试次数未达上限
            if (exception instanceof RecordTooLargeException && retryCounter.get() < MAX_RETRY_TIMES) {
                int currentRetry = retryCounter.incrementAndGet();
                System.out.printf("消息发送失败,开始第%d次重试%n", currentRetry);
                
                // 退避等待,高并发场景建议将等待+重试任务提交到业务线程池执行,避免阻塞Kafka IO线程
                try {
                    Thread.sleep(RETRY_BACKOFF_MS);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    handlePermanentFailure(record, e);
                    return;
                }

                // 此处可添加自定义逻辑:比如压缩消息、拆分消息、调整请求大小配置后再重试
                doSend(record, retryCounter);
                return;
            }

            // 重试耗尽或非重试异常,执行兜底处理
            handlePermanentFailure(record, exception);
        });
    }

    private void handlePermanentFailure(ProducerRecord<String, byte[]> record, Exception ex) {
        // 自定义永久失败逻辑:发送到DLT死信队列、写入数据库存储等
        System.out.printf("消息永久失败,异常信息:%s,准备执行兜底存储%n", ex.getMessage());
    }
}
注意事项
  • 重试逻辑中建议加入退避间隔,不要无间隔连续重试,避免造成服务负载突增
  • 高并发场景下不要在Kafka生产者的IO回调线程中执行阻塞操作,重试任务建议转交业务线程池处理
  • 如果重试时没有调整消息大小或配置的逻辑,不建议对RecordTooLargeException做重试,无意义的重试只会浪费系统资源

内容的提问来源于stack exchange,提问作者user3575226

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 09:27:15