Spring Boot中如何为Kafka内部变更日志主题配置专属Producer参数
我的Kafka Streams拓扑生成了如KSTREAM-AGGREGATE-STATE-STORE-0000000031的内部状态存储主题,同时会自动创建对应的变更日志主题<app-id>KSTREAM-AGGREGATE-STATE-STORE-0000000031-changelog。
拓扑处理器片段如下:
<...> Processor: KSTREAM-FLATMAPVALUES-0000000022 (stores: []) --> KSTREAM-AGGREGATE-0000000032, KSTREAM-FLATMAP-0000000027, KSTREAM-MAP-0000000023, KSTREAM-MAP-0000000025, KSTREAM-MAP-0000000029 <-- KSTREAM-TRANSFORMVALUES-0000000017 Processor: KSTREAM-AGGREGATE-0000000032 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000031]) --> KTABLE-TOSTREAM-0000000033 <-- KSTREAM-FLATMAPVALUES-0000000022 Processor: KTABLE-TOSTREAM-0000000033 (stores: []) --> KSTREAM-PEEK-0000000034 <-- KSTREAM-AGGREGATE-0000000032 <...>
拓扑代码定义如下(BusObjKey和BusObj均为Avro对象并配有对应Serde,TransformBusObj提供聚合及后续映射的业务逻辑):
<...> KStream<BusObjKey, BusObj> busObjStream = otherBusObjStream .groupByKey() .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5))) .aggregate(BusObj::new, TransformBusObj::aggregate, Materialized.with(busObjKeySerde, busObjSerde)) .toStream() .map(TransformBusObj::map); <...>
我希望控制该变更日志主题对应的Producer的配置,特别是开启snappy压缩(例如设置config.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy")),且不影响其他Producer的配置,请问在Spring Boot中该如何实现?
要单独为这个聚合操作的变更日志主题配置snappy压缩,且不影响其他Producer,只需通过Materialized类的withLoggingEnabled方法,为当前状态存储的变更日志指定专属Producer配置即可,具体步骤如下:
创建专属的变更日志配置Map
定义只针对该变更日志的Producer参数,仅设置压缩类型为snappy:import org.apache.kafka.clients.producer.ProducerConfig; import java.util.HashMap; import java.util.Map; Map<String, Object> changelogConfig = new HashMap<>(); changelogConfig.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy");修改
Materialized配置
在聚合操作的Materialized设置中,替换原有的Materialized.with()调用,通过withLoggingEnabled()传入上述配置,同时保留原有的Serde设置:KStream<BusObjKey, BusObj> busObjStream = otherBusObjStream .groupByKey() .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5))) .aggregate(BusObj::new, TransformBusObj::aggregate, Materialized.<BusObjKey, BusObj, WindowStore<Bytes, byte[]>>with(busObjKeySerde, busObjSerde) .withLoggingEnabled(changelogConfig)) // 传入专属变更日志配置 .toStream() .map(TransformBusObj::map);
原理说明
Kafka Streams的Materialized类允许针对单个状态存储的变更日志(changelog)设置独立的Producer配置,通过withLoggingEnabled传入的配置只会作用于当前聚合操作生成的变更日志主题,不会覆盖全局的Kafka Streams配置,也不会影响其他Producer或其他状态存储的行为,完全符合精准控制单个变更日志Producer的需求。
内容的提问来源于stack exchange,提问作者Thomas

