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

无法自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 15:00:17