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

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捕获到反序列化异常,不会将消息传递给监听器。

解决方案

有两种可行方案:

  1. 修改消费者配置,指定默认的value类型
    在application.yaml的kafka consumer properties中添加以下配置(替换为OrderCreated类的实际包路径):
spring.json.value.default.type: com.your.package.OrderCreated

这样即使消息没有__TypeId__头,消费者也会默认用指定类型反序列化JSON内容。

  1. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 17:42:01