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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 17:51:22