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 {
问题原因及修复方案
核心问题
- 配置冲突:容器工厂已开启
batchListener=true,但application.yml中设置listener.type: single,强制使用单条消息模式,导致批量读取的消息仅传递第一条。 - 方法参数不匹配:当前监听器方法接收单个
Flight对象,而非批量集合类型,无法承载批量消息。
修复步骤
- 修改
application.yml中的监听器类型为batch:
spring: kafka: listener: type: batch # 替换原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 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
相关产品推荐
相关产品推荐

