无法自定义KafkaStreams GlobalStateRestoreListener 日志无自定义输出
Kafka Streams自定义StateRestoreListener日志不生效问题
我配置了自定义的StateRestoreListener,但删除本地状态并重启应用让Kafka Streams从changelog topic重建状态时,日志里看不到自定义的日志信息,只显示默认的恢复完成日志。
复现步骤
- 向Kafka Topic写入数据,KStream读取并将状态存储到本地。
- 停止应用,删除本地状态后重启,Kafka Streams从changelog topic重建状态。
- 日志仅显示默认内容:
[-StreamThread-2] o.a.k.s.p.i.StoreChangelogReade stream-thread [foobar-91eae487-939e-439a-bd5f-c918c1f13145-StreamThread-2] Finished restoring changelog foobar-test-avro-leg-changelog-1 to store test-avro-leg with a total number of 66718 records
我的配置类代码
@EnableAutoConfiguration @Slf4j public class SpringKafkaStreamConfig { @Bean public StreamsBuilderFactoryBeanCustomizer streamsBuilderFactoryBeanCustomizer(){ return factoryBean -> { List< StreamsBuilderFactoryBean.Listener > out = factoryBean.getListeners(); factoryBean.setKafkaStreamsCustomizer(new KafkaStreamsCustomizer() { @Override public void customize(KafkaStreams kafkaStreams) { kafkaStreams.setGlobalStateRestoreListener(new StateRestoreListener() { java.util.Date start = null; java.util.Date stop = null; @Override public void onRestoreStart(TopicPartition topicPartition,String storeName,long startingOffset,long endingOffset) { start = Time.from(Instant.now()); log.info("Restarting the building of the following " + "state store: {} " + "starting " + "at offset: {} at the this time: {}", storeName, startingOffset,Time.from(Instant.now())); } @Override public void onBatchRestored(TopicPartition topicPartition,String storeName,long batchEndOffset,long numRestored) { } @Override public void onRestoreEnd(TopicPartition topicPartition, String storeName,long totalRestored) { stop = Time.from(Instant.now()); log.info("State has completed building at this " + "time: {} and restored for the " + "following records: {}", stop,totalRestored); } }); } }); }; } }
问题原因与解决办法
1. 替换自定义器设置方式
当前代码使用setKafkaStreamsCustomizer会覆盖已有的自定义器,导致配置不生效。改用addKafkaStreamsCustomizer添加自定义逻辑:
factoryBean.addKafkaStreamsCustomizer(new KafkaStreamsCustomizer() { @Override public void customize(KafkaStreams kafkaStreams) { // 原有的StateRestoreListener逻辑保持不变 } });
2. 区分全局/本地状态存储
GlobalStateRestoreListener仅对全局状态存储生效,如果你的状态是本地存储(绑定到StreamThread的普通状态存储),需要通过StreamsConfig配置普通StateRestoreListener:
@Bean public StreamsConfigCustomizer streamsConfigCustomizer() { return config -> { config.put(StreamsConfig.STATE_RESTORE_LISTENER_CLASS_CONFIG, CustomStateRestoreListener.class); }; } @Slf4j public class CustomStateRestoreListener implements StateRestoreListener { private Date start; @Override public void onRestoreStart(TopicPartition topicPartition, String storeName, long startingOffset, long endingOffset) { start = new Date(); log.info("开始恢复状态存储: {}, 起始偏移量: {}, 时间: {}", storeName, startingOffset, start); } @Override public void onBatchRestored(TopicPartition topicPartition, String storeName, long batchEndOffset, long numRestored) { // 可选:添加批量恢复日志 } @Override public void onRestoreEnd(TopicPartition topicPartition, String storeName, long totalRestored) { Date stop = new Date(); log.info("状态存储恢复完成: {}, 总恢复记录数: {}, 结束时间: {}", storeName, totalRestored, stop); } }
3. 检查日志级别
确保日志框架配置中,自定义类所在包的日志级别为INFO或更低,避免日志被过滤。例如Logback配置:
<logger name="你的配置类所在包路径" level="INFO"/>
4. 确认本地状态完全删除
重启前彻底删除Kafka Streams的状态目录(默认路径为/tmp/kafka-streams/<你的应用ID>),确保应用触发changelog恢复流程。
内容的提问来源于stack exchange,提问作者user1238097
相关产品推荐
相关产品推荐

