使用Kafka Streams窗口Join时自动生成Avro Schema的原因咨询
KSTREAM-JOINTHIS-0000000125-store-changelog-value的Avro Schema? 这个现象其实是Kafka Streams窗口Join机制和状态存储容错设计的必然结果,我来一步步给你拆解原因:
1. 窗口Join依赖状态存储暂存数据
窗口Join的核心逻辑是要在指定的时间窗口内,匹配两个流中关联的记录。为了做到这一点,Kafka Streams会自动为参与Join的每个流创建窗口状态存储——简单说就是把每个流中还在窗口有效期内的记录暂存起来,等另一个流中对应窗口的记录进来时做匹配。
2. Changelog主题保证状态可恢复
Kafka Streams是分布式流处理框架,必须考虑应用重启、节点故障或分区重新平衡的场景。为了让状态存储能在这些场景下恢复,框架会为每个状态存储自动创建一个changelog主题:所有对状态存储的增、删、改操作都会被同步写入这个主题。当需要恢复状态时,只需要回放这个主题的消息即可。
3. 自动生成的Schema命名规则
你看到的那个长Schema名称,是Kafka Streams针对changelog主题的默认命名逻辑:
KSTREAM-JOINTHIS:标识这是Join操作中用来暂存“待匹配”记录的状态存储(对应Join流程中的其中一个流)0000000125:状态存储的唯一ID,由框架自动生成,用来区分不同的状态存储store-changelog-value:表明这是该状态存储对应的changelog主题中,value部分的Schema
4. 你的Avro Serde触发了Schema自动注册
你代码中使用了SpecificAvroSerde来序列化流中的FactCallProviderMessage,并且配置了Schema Registry地址。当Kafka Streams需要序列化changelog主题的消息时,会复用你为流配置的Serde(因为changelog的数据格式和流中数据格式一致)。而Avro Serde在第一次序列化数据时,会自动将对应的Schema注册到Schema Registry——这就是你看到那个自动命名Schema的来源。
补充个小验证:你可以登录Schema Registry查看这个自动生成的Schema内容,会发现它和你定义的FactCallProviderMessage的Avro Schema完全一致,只是名称被框架替换成了和changelog主题绑定的命名。
你的相关代码参考
Serde<FactCallProviderMessage> specificAvroSerdeForCallProviderMessage = ProcessStreamUtil.getAndRegisterSerde(isKeySerde); KStream<String, FactCallProviderMessage> callProviderMessageKStream = builder.stream( callProviderMessageTopic /* input topic */, Consumed.with(Serdes.String(), specificAvroSerdeForCallProviderMessage)); public static <T extends SpecificRecord> Serde<T> getAndRegisterSerde(boolean isKeySerde) { Serde<T> specificAvroSerde = new SpecificAvroSerde<T>(); specificAvroSerde.configure(Collections.singletonMap( AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, MyConfig.getSchemaRegistryUrl()), isKeySerde); return specificAvroSerde; }
内容的提问来源于stack exchange,提问作者Weili Sun

