Kafka配置25秒超时,FutureRecordMetadata无参get()为何仍60秒超时?
我对Kafka的FutureRecordMetadata类的无参get()方法工作机制有疑问。我的服务中有一个生产消息到topic的方法,希望仅在收到broker的ack后再返回结果。我为生产者配置了以下属性:
kafka.linger-ms=5 kafka.acks=all kafka.enable-idempotence-config=true kafka.delivery-timeout-ms=25000 kafka.max-in-flight=5 kafka.retry-backoff-ms=1000 kafka.request-timeout-ms=20000
服务方法代码如下:
Future<RecordMetadata> result = this.someProducer.send( new ProducerRecord<>( this.someTopic, null, someDTO ) ); try { result.get(); return someDTO.getId(); } catch (ExecutionException e) { throw new NotProducedException("Failed to produce : " + e.getMessage(), e); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new NotProducedException("Interrupted while producing", e); }
我的问题并非如何配置get()使其等待25秒,而是为何实际需要等待60秒,而非我指定的超时时间?
核心原因是你使用了错误的Kafka生产者配置属性名,导致delivery.timeout.ms等关键配置未生效,生产者使用了默认值。
关键问题点
配置属性名格式错误:Kafka生产者的配置属性采用点分隔命名(如
delivery.timeout.ms),而非你使用的短横线分隔(delivery-timeout-ms)。你的配置中大部分属性名不符合Kafka规范,导致这些配置被生产者忽略,转而使用默认值。- 比如
delivery.timeout.ms的默认值正好是60000毫秒(60秒),这就是你看到等待60秒的直接原因。 - 其他错误的属性名包括:
linger-ms(应为linger.ms)、enable-idempotence-config(应为enable.idempotence)、max-in-flight(应为max.in.flight.requests.per.connection)等。
- 比如
无参
get()的行为:无参Future.get()本身没有超时限制,它会一直阻塞直到Kafka生产者完成消息投递(成功收到ack或最终投递失败)。而生产者何时完成投递,完全由其内部的超时逻辑决定——当你的delivery.timeout.ms配置未生效时,生产者会使用默认的60秒超时,因此get()会阻塞到这个超时时间结束。
修正后的配置示例
将配置属性名改为Kafka标准格式:
linger.ms=5 acks=all enable.idempotence=true delivery.timeout.ms=25000 max.in.flight.requests.per.connection=5 retry.backoff.ms=1000 request.timeout.ms=20000
额外说明
当enable.idempotence=true时,生产者会自动设置retries为Integer.MAX_VALUE,但delivery.timeout.ms会限制总投递时间(包括重试和退避等待时间),确保在指定时间内完成或失败。只要配置正确,生产者会在25秒内结束投递逻辑,get()也会随之返回或抛出异常。
内容的提问来源于stack exchange,提问作者Georgi Traykov

