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

Spring Cloud Stream Supplier中Kafka Json Schema序列化异常排查

问题原因

1. 默认序列化器不匹配

Spring Cloud Stream Kafka Binder默认使用ByteArraySerializer作为值序列化器,通过Supplier生产消息时,若未显式配置JSON Schema序列化器,消息会被编码为二进制,Schema Registry只能识别到Byte[]类型的Schema。而直接使用KafkaProducer时,你可能手动指定了JsonSchemaSerializer,因此能正确生成PaymentEvent对应的Schema。

2. 绑定配置未关联Schema Registry参数

Spring Cloud Stream的绑定需要明确配置Schema Registry相关的序列化器参数(包括注册中心地址、目标类类型等)。如果这些配置缺失,Supplier会走默认序列化逻辑,无法生成正确的JSON Schema。

解决方法

1. 配置正确的序列化器与Schema Registry参数

在application.yml(或application.properties)中为Supplier绑定添加以下配置:

YAML 示例:

spring:
  cloud:
    stream:
      kafka:
        bindings:
          payment-out-0: # 替换为你的Supplier绑定名称,默认格式为<函数名>-out-0
            producer:
              value-serializer: io.confluent.kafka.serializers.json.JsonSchemaSerializer
              configuration:
                schema.registry.url: http://你的SchemaRegistry地址:8081
                value.subject.name.strategy: io.confluent.kafka.serializers.subject.RecordNameStrategy
                json.value.type: com.你的包名.PaymentEvent # 替换为PaymentEvent的全类名
      function:
        definition: paymentSupplier # 替换为你的Supplier函数名称

Properties 示例:

spring.cloud.stream.kafka.bindings.payment-out-0.producer.value-serializer=io.confluent.kafka.serializers.json.JsonSchemaSerializer
spring.cloud.stream.kafka.bindings.payment-out-0.producer.configuration.schema.registry.url=http://你的SchemaRegistry地址:8081
spring.cloud.stream.kafka.bindings.payment-out-0.producer.configuration.value.subject.name.strategy=io.confluent.kafka.serializers.subject.RecordNameStrategy
spring.cloud.stream.kafka.bindings.payment-out-0.producer.configuration.json.value.type=com.你的包名.PaymentEvent
spring.cloud.stream.function.definition=paymentSupplier

2. 确认依赖完整性

检查项目依赖是否包含Confluent JSON Schema序列化器和Spring Cloud Stream Kafka Binder:

Maven 依赖:

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-binder-kafka</artifactId>
</dependency>
<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-json-schema-serializer</artifactId>
    <version>${confluent.version}</version> <!-- 使用与Kafka版本兼容的Confluent版本 -->
</dependency>

Gradle 依赖:

implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka'
implementation 'io.confluent:kafka-json-schema-serializer:$confluentVersion'

3. 显式指定Supplier输出类型

在Supplier的Bean定义中明确指定输出类型为PaymentEvent,避免类型擦除导致序列化器无法识别目标类:

@Bean
public Supplier<PaymentEvent> paymentSupplier() {
    return () -> {
        PaymentEvent event = new PaymentEvent();
        event.setId(1L);
        event.setAmount(new BigDecimal("100.00"));
        // 根据业务逻辑填充字段
        return event;
    };
}

4. 同步消费者配置

确保消费者端也配置对应的JSON Schema反序列化器,参数与生产者保持一致:

spring:
  cloud:
    stream:
      kafka:
        bindings:
          payment-in-0: # 消费者绑定名称
            consumer:
              value-deserializer: io.confluent.kafka.serializers.json.JsonSchemaDeserializer
              configuration:
                schema.registry.url: http://你的SchemaRegistry地址:8081
                specific.avro.reader: true
                json.value.type: com.你的包名.PaymentEvent
验证方式
  1. 启动应用后,通过Schema Registry的API或UI检查对应的Subject,确认生成的是PaymentEvent的JSON Schema而非Byte[]类型。
  2. 使用Kafka工具(如kafka-console-consumer.sh)查看消息内容,确认是JSON格式而非二进制。
  3. 启动消费者服务,验证能否正常接收并处理PaymentEvent消息。

内容的提问来源于stack exchange,提问作者Nicholas Irving

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 20:40:10