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

