Spring Cloud Stream消息体嵌入contentType头如何去除?
我完全懂你遇到的糟心事——用Spring Cloud Stream 1.3.2.RELEASE发送String类型消息到Kafka后,不管是用命令行消费者还是Spring Kafka的@KafkaListener,都能看到原本干净的消息体里被硬塞了contentType相关的头信息,导致JSON payload被污染。
这个问题本质是Spring Cloud Stream默认会把消息元数据(比如contentType)嵌入到消息体中,尤其是在没有配置合适消息转换器的情况下。下面给你几个实用的解决办法:
方法一:配置原生编码发送
在生产者的配置文件(application.yml或application.properties)中,给对应的绑定通道设置content-type为application/octet-stream,同时开启原生编码,让消息以原始字节形式发送,避免附加元数据:
YAML配置示例
spring: cloud: stream: bindings: test: # 对应你代码里的channel.test()通道名 destination: test # Kafka主题名 content-type: application/octet-stream producer: use-native-encoding: true
Properties配置示例
spring.cloud.stream.bindings.test.destination=test spring.cloud.stream.bindings.test.content-type=application/octet-stream spring.cloud.stream.bindings.test.producer.use-native-encoding=true
方法二:自定义String消息转换器
如果需要更灵活的控制,可以自定义一个消息转换器,直接跳过元数据嵌入逻辑。创建一个StringMessageConverter的Bean替换默认转换器:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.converter.StringMessageConverter; @Configuration public class StreamConfig { @Bean public StringMessageConverter stringMessageConverter() { return new StringMessageConverter(); } }
这个转换器会直接将String类型的payload以原始字符串形式发送,不会在消息体中添加任何额外头信息。
方法三:直接使用KafkaTemplate发送(绕过Spring Cloud Stream)
如果上面的方法都不符合你的需求,也可以直接用Spring Kafka的KafkaTemplate发送消息,完全绕过Spring Cloud Stream的消息包装逻辑:
import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; @Component public class KafkaProducer { private final KafkaTemplate<String, String> kafkaTemplate; public KafkaProducer(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessage() { kafkaTemplate.send("test", "{\"foo\":\"bar\"}"); } }
验证效果
配置完成后,用Kafka命令行消费者测试:
$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test
就能看到干净的消息体:{"foo":"bar"}
用Spring Kafka的@KafkaListener消费时,也能直接拿到正常的字符串:
@KafkaListener(topics = "test") public void receive(@Payload String message){ log.info("Message payload received: {}", message); }
日志输出会是:2018-05-16 07:16:14.313 INFO 19747 --- [ntainer#0-0-C-1] com.demo.service.Listener : Message payload received: {"foo":"bar"}
原因补充
Spring Cloud Stream 1.x版本中,默认的消息序列化机制会把消息头信息(比如contentType)和payload一起序列化到消息体里,目的是跨消息中间件时保留元数据。但对于Kafka这种原生支持消息头的中间件来说,这种嵌入方式反而会造成污染。通过上面的配置,我们可以让消息以原生形式发送,避免头信息侵入payload。
内容的提问来源于stack exchange,提问作者wltheng

