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

Spring Kafka批量消费异常:仅收到批次中第一条消息

Spring Kafka批量消费异常:反序列化器读取5条消息但监听器仅收到第一条

我们开发了一个带自定义反序列化器的Spring Kafka应用,用@KafkaListener注解接收消息。在自定义反序列化器中添加日志后发现,系统已按批次大小5读取了预期数量的消息,但标注@KafkaListener的方法仅收到该批次中的第一条消息。

Kafka配置类

package com.aa.ctlctr.processor.config;

import java.util.HashMap;
import java.util.Map;

import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.log4j.Logger;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.PropertySource;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

import com.aa.opshub.msgnode.flight.event.json.model.Flight;

@EnableKafka
@Configuration
@PropertySource(value = "classpath:application.yml")
@EnableConfigurationProperties
public class KafkaSourceConfig {
    
    private static Logger logger=Logger.getLogger(KafkaSourceConfig.class);
    
    @Value("${spring.kafka.bootstrap-servers}")
    private String brokerConnect;
    
    @Value("${spring.kafka.consumer.enable-auto-commit}")
    private boolean enableAutocommit;
    
    @Value("${spring.kafka.listener.ack-mode}")
    private String groupIdConfig;    
    
    @Value("${spring.kafka.consumer.properties.max.poll.records:5}")
    private String maxPollRecordConfig;
    
    @Value("${spring.kafka.properties.security.protocol}")
    private String securityProtocol;
      
    @Value("${spring.kafka.properties.sasl.mechanism}")
    private String saslMechanism;
    
    @Value("${spring.kafka.properties.sasl.jaas.config}")
    private String saslJaasConfig;
    
    @Value("${spring.kafka.properties.sasl.login.callback.handler.class}")
    private String saslClientCallbackHandlerClass;
    
    
    @Bean
    public Map<String, Object> consumer_Configs() {
        Map<String, Object> prop = new HashMap<>();
        prop.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerConnect);
        prop.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaCustomDeserializer.class);
        prop.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutocommit);
        prop.put(ConsumerConfig.GROUP_ID_CONFIG, groupIdConfig);
        prop.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecordConfig);
        prop.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
        prop.put(SaslConfigs.SASL_MECHANISM, saslMechanism);
        prop.put(SaslConfigs.SASL_JAAS_CONFIG, saslJaasConfig);
        prop.put(SaslConfigs.SASL_CLIENT_CALLBACK_HANDLER_CLASS, saslClientCallbackHandlerClass);
        prop.put(SaslConfigs.SASL_JAAS_CONFIG, saslJaasConfig);
        prop.put("ssl.engine.factory.class", InsecureSslEngineFactory.class);
        return prop;
    }


    @Bean
    public ConsumerFactory<String, Flight> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumer_Configs(),new StringDeserializer(), 
                  new KafkaCustomDeserializer<>());

    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Flight> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Flight> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setBatchListener(true);
        return factory;
    }   
 
}

application.yml配置

spring:
  kafka:
    #bootstrap-servers: ${kafka.bootstrap.servers}
    bootstrap-servers: <<broker address>>
    properties:
      security:
        protocol: SASL_SSL
      sasl:
        mechanism: OAUTHBEARER
        jaas:
          config: org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required;
        login:
          callback:
            handler:
              class: <<Security call back handler>>
      max.request.size: 750000
      request.timeout.ms: 30000
      linger.ms: 500
      delivery.timeout.ms: 91500
      metadata.max.age.ms: 180000
      connections.max.idle.ms: 60000
    consumer:          
      enable-auto-commit: false
      auto-offset-reset: earliest
      properties:
        max.poll.records: 5
        partition.assignment.strategy: org.apache.kafka.clients.consumer.RoundRobinAssignor
    listener:
      type: single
      ack-mode: batch 

@KafkaListener方法

@KafkaListener(topics = "#{'${my.kafka.conf.topics}'.split(',')}", concurrency = "${my.kafka.conf.concurrency}", clientIdPrefix = "${my.kafka.conf.clientIdPrefix}", groupId = "${my.kafka.conf.groupId}")
public void kafkaListener(final Flight flight,@Header(KafkaHeaders.OFFSET) Long offset,
        @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
        @Header(KafkaHeaders.RECEIVED_TIMESTAMP) Long timestamp) throws JsonMappingException, JsonProcessingException {

问题原因及修复方案

核心问题

  1. 配置冲突:容器工厂已开启batchListener=true,但application.yml中设置listener.type: single,强制使用单条消息模式,导致批量读取的消息仅传递第一条。
  2. 方法参数不匹配:当前监听器方法接收单个Flight对象,而非批量集合类型,无法承载批量消息。

修复步骤

  1. 修改application.yml中的监听器类型为batch:
spring:
  kafka:
    listener:
      type: batch  # 替换原single配置
      ack-mode: batch 
  1. 更新@KafkaListener方法参数,改为接收批量消息集合,同时头部参数对应改为集合类型:
@KafkaListener(topics = "#{'${my.kafka.conf.topics}'.split(',')}", concurrency = "${my.kafka.conf.concurrency}", clientIdPrefix = "${my.kafka.conf.clientIdPrefix}", groupId = "${my.kafka.conf.groupId}")
public void kafkaListener(final List<Flight> flights,
        @Header(KafkaHeaders.OFFSET) List<Long> offsets,
        @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitionIds,
        @Header(KafkaHeaders.RECEIVED_TIMESTAMP) List<Long> timestamps) throws JsonMappingException, JsonProcessingException {
    // 批量消息处理逻辑
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 22:30:50