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

Spring Cloud Stream Kafka Streams拓扑Transform步骤Serde配置异常

问题描述

使用Spring Cloud Stream Kafka Streams Binder 3.2.4、Confluent Platform 7.1.1开发,Key和Value均采用SpecificAvroSerde实现Serde,因主题包含带引用的联合Schema,已配置auto_register_schemas=false。但在拓扑的Transform步骤后,流抛出序列化错误,提示Schema Registry无法找到内部重分区主题对应的Schema,报错信息如下:

Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Subject 'utw-local-095515-t-tipsProcess-repartition-value' not found.; error code: 40401 

补充背景:单个主题中使用多种事件类型,已按Confluent指南禁用自动Schema注册并完成主主题的Schema注册,但Transform后生成的内部重分区主题未关联已注册的Schema。

相关代码:

@Bean
@Autowired
public Consumer<KStream<UnreportedTipsKey, SpecificRecord>> tipsProcess(
        @Value("${kafka.store.tips}") String unreportedTipsStore,
        Serde<UnreportedTipsKey> tipsKeySerde, Serde<SpecificRecord> specificSerde,
        EventPersister eventPersister, @Value("${kafka.topic.tips}") String topic,
        TransformerSupplier<UnreportedTipsKey, SpecificRecord, KeyValue<UnreportedTipsKey, SpecificRecord>> tipsTransformerSupplier,
        @Value("${MDE_DEPLOYMENT_STAGE}") String env) {

    return input -> input.transform(tipsTransformerSupplier)
            .groupByKey(Grouped.with(tipsKeySerde, specificSerde))
            .aggregate(UnreportedTips::new,
                    new UnreportedTipsAggregator(
                            eventPersister, topic, env),
                    Materialized
                            .<UnreportedTipsKey, UnreportedTips, KeyValueStore<Bytes, byte[]>>as(
                                    unreportedTipsStore)
                            .withKeySerde(tipsKeySerde).withValueSerde(
                                    SerdeConfiguration.unreportedTipsSerde()));
}
@Configuration
public class SerdeConfiguration {

    private static final String SCHEMA_REGISTRY_URL_CONFIG = "schema.registry.url";
    private static final String AUTO_REGISTER_SCHEMAS = "auto.register.schemas";

    @Bean
    public Serde<UnreportedTipsKey> tipsKeySerde(
            @Value("${spring.kafka.properties.schema.registry.url}") String schemaRegistryUrl) {

        final Map<String, Object> serdeConfig = new HashMap<>();
        serdeConfig.put(SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        serdeConfig.put(AUTO_REGISTER_SCHEMAS, false);

        final Serde<UnreportedTipsKey> keySpecificRecordSerde = new SpecificAvroSerde<>();
        keySpecificRecordSerde.configure(serdeConfig, true);

        return keySpecificRecordSerde;
    }

    @Bean
    public Serde<SpecificRecord> specificRecordValueSerde(
            @Value("${spring.kafka.properties.schema.registry.url}") String schemaRegistryUrl) {

        final Map<String, Object> serdeConfig = new HashMap<>();
        serdeConfig.put(SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        serdeConfig.put(AUTO_REGISTER_SCHEMAS, false);

        final Serde<SpecificRecord> specificRecordValueSerde = new SpecificAvroSerde<>();
        specificRecordValueSerde.configure(serdeConfig, false);

        return specificRecordValueSerde;

    }

    public static Serde<UnreportedTips> unreportedTipsSerde() {
        Serdes.String();
        return new JsonSerde<>(new ObjectMapper().constructType(UnreportedTips.class));
    }
}

已尝试的无效方案:

  • 自定义忽略Schema Registry的Serializer,配置spring.cloud.stream.kafka.streams.bindings.tipsProcess.producer.keySerde和spring.cloud.stream.kafka.streams.bindings.tipsProcess.producer.valueSerde,但配置被忽略,仍使用KafkaAvroSerde。
  • 更新Serde配置使用RecordNameStrategy,无效,因为Transform步骤使用了不同的Serde。

困惑:官方文档提到的Value serdes推断规则不适用于当前拓扑(函数未声明输出流),不清楚如何配置中间主题的Serde。

解决方案

1. 全局配置Kafka Streams默认Serde

在配置文件中指定全局默认的Key和Value Serde,让所有内部主题(重分区、changelog等)使用你自定义的Serde:

# 直接引用Serde Bean的名称
spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=tipsKeySerde
spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=specificRecordValueSerde

或者使用Serde的全限定类名:

spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=com.your.package.SerdeConfiguration$TipsKeySerde
spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=com.your.package.SerdeConfiguration$SpecificRecordValueSerde

该配置会覆盖Kafka Streams的默认Serde,确保内部主题序列化时使用你配置好的SpecificAvroSerde。

2. 为Serde配置Subject名称策略

修改Serde配置,添加RecordNameStrategy,让序列化时根据记录的全限定类名查找Schema,而非内部主题名称:

@Bean
public Serde<SpecificRecord> specificRecordValueSerde(
        @Value("${spring.kafka.properties.schema.registry.url}") String schemaRegistryUrl) {

    final Map<String, Object> serdeConfig = new HashMap<>();
    serdeConfig.put(SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
    serdeConfig.put(AUTO_REGISTER_SCHEMAS, false);
    // 设置Subject策略为记录名称,避免依赖内部主题名称
    serdeConfig.put("value.subject.name.strategy", io.confluent.kafka.serializers.subject.RecordNameStrategy.class.getName());

    final Serde<SpecificRecord> specificRecordValueSerde = new SpecificAvroSerde<>();
    specificRecordValueSerde.configure(serdeConfig, false);

    return specificRecordValueSerde;
}

同时对Key Serde做相同配置(如果需要),这样重分区主题在序列化时会使用记录类名对应的Schema,而不是内部主题的Subject。

3. 手动注册内部主题Schema

如果必须使用默认的TopicNameStrategy,可以提前在Schema Registry中注册内部重分区主题的Schema。内部重分区主题名称格式为<application-id>-<function-name>-repartition,即报错中的utw-local-095515-t-tipsProcess-repartition,需要注册Subject utw-local-095515-t-tipsProcess-repartition-value,Schema内容与主主题一致即可。

4. 优化拓扑避免不必要的重分区

检查transform步骤是否修改了Key,如果未修改Key,groupByKey不需要触发重分区。可以通过以下方式优化:

return input -> input.transform(tipsTransformerSupplier)
        .groupByKey(Grouped.with(tipsKeySerde, specificSerde))
        .aggregate(UnreportedTips::new,
                new UnreportedTipsAggregator(eventPersister, topic, env),
                Materialized.as(unreportedTipsStore)
                        .withKeySerde(tipsKeySerde)
                        .withValueSerde(SerdeConfiguration.unreportedTipsSerde()))
        .optimize(); // 触发拓扑优化,跳过不必要的重分区

Kafka Streams会自动判断是否需要重分区,若Key未变更则跳过内部重分区主题的生成,从根源解决问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 03:01:11