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

Python与Spring Boot Kafka交互报错:头部无类型信息且无默认类型

问题

正在构建不受编程语言或框架限制的纯微服务,采用Kafka搭建事件管道。当前从Python应用发送JSON格式消息:

producer = KafkaProducer(bootstrap_servers=broker)
# req 是字典 {'a':1, 'c':-2}
json_message = json.dumps(req).encode('utf-8')
print("--------------------------")
print(json_message)
print("--------------------------")
headers = [
     ('content-type', 'application/json'),
     ('type', 'com.mua.cloud.testm.models.events'),
]
producer.send(topic+"-m", json_message,headers=headers)

尝试在Spring Boot应用中消费该消息:

@Service
@Log4j2
public class TestHandler {

   @KafkaListener(topics = Constants.TOPIC_PREFIX + "-" + "python-trigger-m")
   @SneakyThrows
   public void listener(@Header(KafkaHeaders.CORRELATION_ID) byte[] corrId, TestEvent event) {
       System.out.println(corrId);
       System.out.println(event);
       System.out.println("----------------");
   }
}

已修改配置:

spring:
 …
 kafka:
   …
   consumer:
     value-deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer
     value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

并注册Bean:

@Configuration
public class JsonMessageConverterConfig {
   @Bean
   public JsonMessageConverter jsonMessageConverter() {
       return new ByteArrayJsonMessageConverter();
   }
}

但反复报错:

TestMicro  ERROR    2023-07-18 19:08:08,922   [consumer-0-C-1] [o.s.k.l.KafkaMessageListenerContainer] - Message:  Consumer exception 
java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer
    at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:151)
    ...
Caused by: java.lang.IllegalStateException: No type information in headers and no default type provided
    at org.springframework.util.Assert.state(Assert.java:76)
    ...

也曾尝试添加配置spring.kafka.producer.properties.spring.json.add.type.headers=false,还尝试修改consumer配置并设置TYPE_MAPPINGS,但均未解决问题。已知:Spring Boot向Python发送消息正常;修改配置可能破坏现有功能。


解决方案

报错核心原因:Spring Kafka反序列化时,无法确定目标类型——Python发送的消息没有携带Spring Kafka默认识别的__TypeId__头,且未配置默认转换类型。以下是三种不破坏现有功能的解决方式:

方式1:手动反序列化(最简方案)

修改监听方法,先接收byte数组,再用ObjectMapper手动反序列化为TestEvent:

@Service
@Log4j2
public class TestHandler {

    @Autowired
    private ObjectMapper objectMapper;

    @KafkaListener(topics = Constants.TOPIC_PREFIX + "-" + "python-trigger-m")
    @SneakyThrows
    public void listener(@Header(KafkaHeaders.CORRELATION_ID) byte[] corrId, 
                         @Payload byte[] payload) {
        TestEvent event = objectMapper.readValue(payload, TestEvent.class);
        System.out.println(corrId);
        System.out.println(event);
        System.out.println("----------------");
    }
}

方式2:配置专用消费者工厂(隔离现有配置)

创建仅用于消费Python消息的消费者容器工厂,通过自定义消息转换器从Python的type头识别目标类型:

@Configuration
public class PythonKafkaConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public ConsumerFactory<String, TestEvent> pythonConsumerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "python-event-group");
        configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);
        return new DefaultKafkaConsumerFactory<>(configProps, new StringDeserializer(), 
                new ByteArrayDeserializer());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, TestEvent> pythonKafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, TestEvent> factory = 
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(pythonConsumerFactory());
        factory.setMessageConverter(new ByteArrayJsonMessageConverter() {
            @Override
            protected Class<?> getType(MessageHeaders headers) {
                byte[] typeHeaderBytes = (byte[]) headers.get("type");
                if (typeHeaderBytes != null) {
                    String type = new String(typeHeaderBytes);
                    if ("com.mua.cloud.testm.models.events".equals(type)) {
                        return TestEvent.class;
                    }
                }
                return super.getType(headers);
            }
        });
        return factory;
    }
}

修改监听方法指定使用该工厂:

@KafkaListener(topics = Constants.TOPIC_PREFIX + "-" + "python-trigger-m", containerFactory = "pythonKafkaListenerContainerFactory")
@SneakyThrows
public void listener(@Header(KafkaHeaders.CORRELATION_ID) byte[] corrId, TestEvent event) {
    System.out.println(corrId);
    System.out.println(event);
    System.out.println("----------------");
}

方式3:配置ErrorHandlingDeserializer并指定默认类型

在现有消费者配置中添加错误处理反序列化器,同时指定默认转换类型:

spring:
  kafka:
    consumer:
      value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
      properties:
        spring.deserializer.value.delegate.class: org.apache.kafka.common.serialization.ByteArrayDeserializer
        spring.json.value.default.type: com.mua.cloud.testm.models.events.TestEvent # 替换为你的TestEvent全类名

此方式会让Spring Kafka默认将byte数组消息反序列化为指定类型,同时处理序列化错误,不影响现有Spring Boot之间的消息交互。


内容的提问来源于stack exchange,提问作者Maifee Ul Asad

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:50:41