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
相关产品推荐
相关产品推荐

