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

Spring Cloud Stream Kafka生产消息Payload为字节数组而非JSON的问题

解决方案:Spring Cloud Stream生成JSON格式Payload而非字节数组

核心问题分析

你遇到的问题本质是Spring Cloud Stream结合Avro时,默认触发了Avro二进制序列化逻辑,即使配置了contentType: application/json,也因依赖冲突、配置优先级或自动配置逻辑导致未生效。以下是分步解决方法:


1. 清理依赖,解决版本冲突

你的依赖中spring-cloud-stream-schema 2.2.1.RELEASE是旧版本,与spring-cloud-stream 3.2.2不兼容,会干扰消息转换逻辑。

  • 移除spring-cloud-stream-schema依赖
  • 添加Jackson Avro序列化模块(用于将Avro对象转为JSON)
<!-- Maven依赖示例 -->
<dependency>
    <groupId>com.fasterxml.jackson.dataformat</groupId>
    <artifactId>jackson-dataformat-avro</artifactId>
    <version>2.15.2</version> <!-- 版本需与Spring Boot版本匹配,如2.7.x对应2.15.x -->
</dependency>

2. 调整Spring Cloud Stream配置

在application.yml中明确配置输出绑定的消息转换规则,强制使用JSON序列化:

spring:
  cloud:
    stream:
      bindings:
        # 替换为你的输出绑定名称
        your-output-binding:
          destination: target-topic
          contentType: application/json
          producer:
            # 禁用原生编码,让Spring Cloud Stream通过内置转换器处理消息
            use-native-encoding: false
      kafka:
        binder:
          configuration:
            # 指定JSON序列化器,覆盖Avro序列化器
            value.serializer: org.springframework.kafka.support.serializer.JsonSerializer
            key.serializer: org.apache.kafka.common.serialization.StringSerializer

3. 修正响应式函数式代码

确保Avro对象被序列化为JSON字符串(可手动控制或依赖Jackson自动转换):

import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import java.util.function.Function;

public class MessageHandler {

    private final ObjectMapper objectMapper;

    public MessageHandler(ObjectMapper objectMapper) {
        this.objectMapper = objectMapper;
    }

    public Function<Flux<Message<Object>>, Flux<Message<String>>> handler() {
        return inputFlux -> inputFlux.flatMap(message -> {
            FirstRankPaymentAgreed avroPayload = (FirstRankPaymentAgreed) message.getPayload();
            try {
                // 将Avro对象序列化为JSON字符串
                String jsonPayload = objectMapper.writeValueAsString(avroPayload);
                return Mono.just(MessageBuilder.withPayload(jsonPayload)
                        .copyHeaders(message.getHeaders())
                        .setHeader("contentType", "application/json")
                        .build());
            } catch (Exception e) {
                return Mono.error(e);
            }
        });
    }
}

4. 排查额外干扰配置

  • 移除所有与Avro Schema Registry相关的强制序列化配置(如spring.cloud.stream.schema.registry.client.endpoint),除非你需要同时支持Avro和JSON双格式
  • 检查是否有自定义的MessageConverter或Serializer bean覆盖了默认的JSON转换逻辑

内容的提问来源于stack exchange,提问作者Abdullah LAKHDAR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 17:20:42