Kafka重试配置未按预期生效,Broker不可用时如何配置重试?
问题描述
我使用Kafka 2.8.1版本搭建了包含1个Broker、1个Topic、1个Partition的基础环境,使用以下代码发送消息:
public static void main(String[] args) throws InterruptedException { //1. Create producer configuration information Properties properties = new Properties(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:19092"); properties.put("max.block.ms","5000"); properties.put("retries","3"); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); //2. Create producer object KafkaProducer<String,String> producer = new KafkaProducer<String, String>(properties); final CountDownLatch countDownLatch = new CountDownLatch(1); //3. Send data producer.send(new ProducerRecord<String, String>("test","CallBack--1"), new Callback() { public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception == null){ System.out.println(metadata.partition()+"=="+metadata.offset()); countDownLatch.countDown(); }else { System.out.println("Error!!!"); exception.printStackTrace(); } } }); countDownLatch.await(60, TimeUnit.SECONDS); //4. Close resources producer.close(); }
当Broker与Zookeeper均处于宕机状态时,我预期日志中会出现3次“Error!!!”,但仅抛出一次org.apache.kafka.common.errors.TimeoutException: Topic test not present in metadata after 5000 ms错误后无后续重试操作,请问如何配置才能在Broker不可用时实现消息发送重试?
问题原因
出现该问题的核心原因有两点:
- 元数据获取超时提前终止流程:你设置的
max.block.ms=5000(仅5秒),当Broker和ZK宕机时,Producer无法获取目标Topic的元数据,send()方法会阻塞等待元数据直到超时,此时直接抛出TimeoutException,Producer还未进入消息发送的重试阶段,因此retries配置未生效。 - 异常类型不在默认重试范围内:元数据获取超时的
TimeoutException默认不属于Producer自动重试的异常类别,不会触发retries配置的重试逻辑。
解决方案
通过调整以下Producer配置,可实现Broker不可用时的消息发送重试:
- 调整
max.block.ms:设置为与delivery.timeout.ms一致的值,避免因元数据获取超时提前终止发送流程,给Producer足够时间尝试重试。 - 配置
delivery.timeout.ms:控制消息从发送到最终判定失败的总时长,需确保该值大于request.timeout.ms+retry.backoff.ms*retries,保证重试流程能完整执行。 - 明确
retries次数:保留或明确设置重试次数,配合delivery.timeout.ms使用。 - 设置
retry.backoff.ms:定义两次重试的间隔时间,避免频繁重试带来的资源消耗。
修改后的Producer配置代码示例:
public static void main(String[] args) throws InterruptedException { //1. Create producer configuration information Properties properties = new Properties(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:19092"); // 调整max.block.ms,与delivery.timeout.ms一致 properties.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, "60000"); // 设置总超时时间,控制消息从发送到失败的最长时长 properties.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, "60000"); // 明确重试次数 properties.put(ProducerConfig.RETRIES_CONFIG, "3"); // 设置重试间隔,避免频繁重试 properties.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, "1000"); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer"); //2. Create producer object KafkaProducer<String,String> producer = new KafkaProducer<String, String>(properties); final CountDownLatch countDownLatch = new CountDownLatch(1); //3. Send data producer.send(new ProducerRecord<String, String>("test","CallBack--1"), new Callback() { public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception == null){ System.out.println(metadata.partition()+"=="+metadata.offset()); countDownLatch.countDown(); }else { System.out.println("Error!!!"); exception.printStackTrace(); } } }); countDownLatch.await(60, TimeUnit.SECONDS); //4. Close resources producer.close(); }
配置说明
delivery.timeout.ms:默认值为120000ms(2分钟),这里设置为60000ms(1分钟),确保在Broker恢复前有足够时间重试。retry.backoff.ms:默认值为100ms,这里设置为1000ms(1秒),减少重试频率,避免无效请求浪费资源。max.block.ms:设置为与delivery.timeout.ms相同的值,确保send()方法不会因元数据获取超时提前抛出异常,让重试逻辑有机会执行。
内容的提问来源于stack exchange,提问作者Inessa
相关产品推荐
相关产品推荐

