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

Spring Cloud Streams自定义Avro Serdes序列化失败问题排查

问题分析

这个序列化异常的核心原因有两点:

  1. GenericAvroSerde未配置Schema Registry参数:GenericAvroSerde依赖Schema Registry获取/注册Avro schema,但配置中仅指定了Serde类型,未提供Schema Registry地址等必要参数,导致Serde未完成初始化就被调用。
  2. 代码与配置的Serde不匹配:代码中使用自定义CustomSerdes.TestRecord()处理TestRecord类型,但application.yml中输入绑定配置的是GenericAvroSerde,两者可能存在类型不兼容(比如GenericAvroSerde反序列化出GenericRecord,但代码期望TestRecord),或者自定义Serde未正确处理Schema Registry配置。
解决方案

1. 完善Schema Registry全局配置

在application.yml中添加Schema Registry的全局配置,确保所有依赖Avro Serde的组件能获取到配置:

spring:
  kafka:
    streams:
      application-id: test-app
      properties:
        schema.registry.url: http://你的SchemaRegistry地址:8081
        # 可选:需自动注册schema时添加
        auto.register.schemas: true
  cloud:
    function:
      definition: process
    stream:
      kafka:
        streams:
          bindings:
            process-in-0:
              consumer:
                key-serde: org.apache.kafka.common.serialization.Serdes$StringSerde
                value-serde: io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
            # 第二个输入绑定同步配置,两个输入Topic结构一致
            process-in-1:
              consumer:
                key-serde: org.apache.kafka.common.serialization.Serdes$StringSerde
                value-serde: io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
            # 输出绑定必须配置对应Serde,避免默认Serde引发异常
            process-out-0:
              producer:
                key-serde: org.apache.kafka.common.serialization.Serdes$StringSerde
                value-serde: com.yourpackage.CustomSerdes$ModifiedTestPayloadSerde

2. 统一Serde使用逻辑

情况A:TestRecord是Avro生成的特定类

如果TestRecord是通过Avro schema生成的Java类,直接使用SpecificAvroSerde替代自定义Serde:

public class CustomSerdes {
    public static Serde<TestRecord> TestRecord() {
        SpecificAvroSerde<TestRecord> serde = new SpecificAvroSerde<>();
        Map<String, String> config = new HashMap<>();
        config.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://你的SchemaRegistry地址:8081");
        serde.configure(config, false); // false表示作为value Serde,true为key Serde
        return serde;
    }

    // 为ModifiedTestPayload实现对应Serde示例
    public static Serde<ModifiedTestPayload> ModifiedTestPayload() {
        SpecificAvroSerde<ModifiedTestPayload> serde = new SpecificAvroSerde<>();
        Map<String, String> config = new HashMap<>();
        config.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://你的SchemaRegistry地址:8081");
        serde.configure(config, false);
        return serde;
    }
}

情况B:坚持使用自定义Serde

若必须使用自定义TestRecordSerde,需确保Serializer和Deserializer正确处理Schema Registry配置,在Serde的configure方法中传递参数:

public class CustomSerdes {
    public static Serde<TestRecord> TestRecord() {
        return new TestRecordSerde();
    }

    public static class TestRecordSerde extends Serdes.WrapperSerde<TestRecord> {
        public TestRecordSerde() {
            super(new TestRecordSerializer(), new TestRecordDeserializer());
        }

        @Override
        public void configure(Map<String, ?> configs, boolean isKey) {
            super.configure(configs, isKey);
            // 将配置传递给Serializer和Deserializer
            serializer.configure(configs, isKey);
            deserializer.configure(configs, isKey);
        }
    }

    // TestRecordSerializer实现configure方法初始化SchemaRegistryClient
    public static class TestRecordSerializer implements Serializer<TestRecord> {
        private SchemaRegistryClient client;

        @Override
        public void configure(Map<String, ?> configs, boolean isKey) {
            String registryUrl = (String) configs.get(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG);
            this.client = new CachedSchemaRegistryClient(registryUrl, 1000);
            // 其他初始化逻辑
        }

        @Override
        public byte[] serialize(String topic, TestRecord data) {
            // 使用client完成schema注册和序列化
            // ...
        }
    }

    // TestRecordDeserializer同理实现configure方法
    public static class TestRecordDeserializer implements Deserializer<TestRecord> {
        private SchemaRegistryClient client;

        @Override
        public void configure(Map<String, ?> configs, boolean isKey) {
            String registryUrl = (String) configs.get(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG);
            this.client = new CachedSchemaRegistryClient(registryUrl, 1000);
            // 其他初始化逻辑
        }

        @Override
        public TestRecord deserialize(String topic, byte[] data) {
            // 使用client完成schema拉取和反序列化
            // ...
        }
    }
}

3. 确保StreamJoined的Serde与输入匹配

代码中StreamJoined.with指定的Serde需和输入绑定配置的Serde兼容:
如果输入用GenericAvroSerde,则代码应处理GenericRecord类型;或者将输入配置改为SpecificAvroSerde,反序列化为TestRecord,避免类型不匹配引发异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 08:14:52