Spring Kafka中linger.ms配置未生效回退默认值问题排查
Kafka生产者自定义批量参数未生效问题排查
业务预期
基于Kafka实现批量消息处理逻辑:1分钟内累计有100条请求到达时,不立即逐条发布消息,将100条请求攒为整批后统一发送到Kafka Topic。
实际异常现象
配置完成后批量发送逻辑未生效:消息一旦触发发送就会立即发布到Topic,同时被消费者接收,未实现攒批效果。
配置中已设置linger.ms = 60000,按照该参数的官方语义,生产者会至少等待指定时长再发送消息,即使发送线程提前空闲、批次大小未达到阈值也不会提前发送。但实际运行中消息完全没有等待60000ms,也未等到batch.size配置的阈值就立即发送。
排查Spring Kafka运行日志发现,自定义的ProducerConfig参数未被正确加载,全部回退为默认值,日志打印的生产者参数如下:
ProducerConfig values: acks = -1 batch.size = 16384 bootstrap.servers = [localhost:9092] buffer.memory = 33554432 client.dns.lookup = use_all_dns_ips client.id = producer-1 compression.type = none connections.max.idle.ms = 540000 delivery.timeout.ms = 120000 enable.idempotence = true interceptor.classes = [] key.serializer = class org.apache.kafka.common.serialization.StringSerializer linger.ms = 0
现有配置代码
生产者配置
public class KafkaProducerConfig { @Bean public Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.ACKS_CONFIG, 1); props.put(ProducerConfig.RETRIES_CONFIG, 0); props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 1); props.put(ProducerConfig.LINGER_MS_CONFIG, 60000); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 100000); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, DefaultPartitioner.class); return props; } @Bean public ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } }
消费者配置
public class KafkaConfig { ConsumerFactory<String, String> kafkaConsumerFactory(Boolean autoCommit) { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, autoCommit); props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 1000); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 20000); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 5); props.put(ConsumerConfig.GROUP_ID_CONFIG, "batch"); return new DefaultKafkaConsumerFactory<>(props); } @Bean("kafkaListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(kafkaConsumerFactory(Boolean.TRUE)); factory.setBatchListener(true); factory.setConcurrency(1); return factory; } }
生产端发送逻辑
@Autowired private KafkaTemplate<String, String> kafka; kafka.send("batch-test", message);
消费端监听逻辑
@KafkaListener(id = "testGroup", topics = {"batch-test"}) public void test(@Payload List<String> messages, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions, @Header(KafkaHeaders.OFFSET) List<Long> offsets) { for (int i = 0; i < messages.size(); i++) { System.out.println(messages.get(i) + partitions.get(i) + "-" + offsets.get(i) + ""); } }
排查目标
定位自定义配置的linger.ms、batch.size等生产者参数未生效、全部回退为默认值的问题原因。
内容的提问来源于stack exchange,提问作者Swastik
相关产品推荐
相关产品推荐

