Spring Kafka监听器无法消费CLI消息但可接收KafkaTemplate消息
问题:CLI发送Kafka消息无法被监听器接收,KafkaTemplate发送正常
按照序列化/反序列化相关文档操作时,出现以下问题:用CLI发送消息到order.created主题时,监听器无法接收;但用Spring的KafkaTemplate发送消息时,监听器能正常处理。
CLI发送消息操作
docker exec -it broker1 kafka-console-producer --bootstrap-server localhost:9092 --topic order.created --property parse.key=true --property key.separator=@ >"123"@{"orderId":"bfd9ebf3-0218-49dd-bc3e-dcd172ec3647"}
KafkaTemplate发送代码
@Autowired KafkaTemplate<String, Object> kafkaTemplate; @GetMapping("/test") public String test(){ kafkaTemplate.send("order.created", "123", new OrderCreated()); return "sent"; }
Spring配置文件(application.yaml)
spring: main: allow-bean-definition-overriding: true banner-mode: 'off' profiles: active: default application: name: kafka-experimental jackson.default-property-inclusion: NON_NULL kafka: admin: fail-fast: true consumer: group-id: order-group # 文档说明: # 当反序列化器无法反序列化消息时,Spring无法处理该问题,因为它发生在poll()返回之前。 # 为解决此问题,引入了ErrorHandlingDeserializer,它委托给实际的反序列化器(键或值)。 # 如果委托反序列化失败,ErrorHandlingDeserializer返回null值,并在头中包含DeserializationException(包含原因和原始字节)。 # 如果ConsumerRecord的键或值包含DeserializationException头,容器的ErrorHandler会处理失败的记录,不会传递给监听器。 key-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer auto-startup: true max-poll-records: 500 properties: acks: all allow.auto.create.topics: true partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor # 键的委托反序列化器 spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer # 值的委托反序列化器 spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer # spring.json.value.default.type: # spring.json.use.type.headers: true spring.json.trusted.packages: "*" security.protocol: "PLAINTEXT" producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer bootstrap-servers: - http://localhost:9092
监听器代码
@KafkaListener( clientIdPrefix = "OrderCreatedHandlerReceiver", id = "orderConsumerClient", topics = "order.created", groupId = "order-group", containerFactory = "kafkaListenerContainerFactory") public void listen( @Header(KafkaHeaders.RECEIVED_PARTITION) Integer partition, @Header(KafkaHeaders.RECEIVED_KEY) String key, @Payload OrderCreated payload, @Headers Map<String, Object> headers) { log.info("Received Kafka message: payload={} with key={} on partition={}", payload, key, partition); }
Broker配置(docker-compose.yaml)
version: '3.8' name: "kafka-experimental-cluster" services: zookeeper: image: confluentinc/cp-zookeeper:6.1.1 hostname: zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 broker1: image: confluentinc/cp-kafka:6.1.1 hostname: broker1 container_name: broker1 depends_on: - zookeeper ports: - "29092:29092" - "9092:9092" - "9101:9101" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker1:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 # 单broker集群设为1,用于跟踪消费者偏移量;若broker宕机会丢失偏移信息,生产环境建议大于1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 # 副本和leader共2个节点时需确认消息接收 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 # 开发环境设1,生产环境设3 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 3000 # 新消费者加入组时的初始重平衡延迟,减少频繁重平衡次数;测试环境可设0,生产环境设3000 AUTO_CREATE_TOPICS: true KAFKA_COMPRESSION_TYPE: producer KAFKA_JMX_PORT: 9101 KAFKA_JMX_HOSTNAME: localhost
问题分析与解决方案
核心原因
使用KafkaTemplate发送消息时,JsonSerializer会自动添加__TypeId__消息头,明确告诉消费者该消息对应的Java类型;但CLI发送的消息没有这个头信息,而当前消费者配置未指定默认的反序列化类型,导致JsonDeserializer无法识别要将JSON字符串反序列化为OrderCreated对象。最终ErrorHandlingDeserializer捕获到反序列化异常,不会将消息传递给监听器。
解决方案
有两种可行方案:
- 修改消费者配置,指定默认的value类型
在application.yaml的kafka consumer properties中添加以下配置(替换为OrderCreated类的实际包路径):
spring.json.value.default.type: com.your.package.OrderCreated
这样即使消息没有__TypeId__头,消费者也会默认用指定类型反序列化JSON内容。
- CLI发送时手动添加类型头
使用kafka-console-producer发送消息时,通过--property指定添加__TypeId__头(替换为实际包路径):
docker exec -it broker1 kafka-console-producer --bootstrap-server localhost:9092 --topic order.created --property parse.key=true --property key.separator=@ --property headers=__TypeId__=com.your.package.OrderCreated
然后输入消息内容:
"123"@{"orderId":"bfd9ebf3-0218-49dd-bc3e-dcd172ec3647"}
关于ErrorHandlingDeserializer的验证说明
你提到发送new Student()时监听器不触发但能在topic中看到消息,这符合预期:Student类型消息的__TypeId__头是Student的类路径,与监听器期望的OrderCreated不匹配,反序列化失败后ErrorHandlingDeserializer会将异常放入消息头,容器错误处理器会处理该消息,不会传递给监听器。
内容的提问来源于stack exchange,提问作者Wrapper
相关产品推荐
相关产品推荐

