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

Kafka Streams手动提交时repartition topic记录未清理问题咨询

Kafka Streams手动提交时Repartition Topic无法自动清理

问题描述

我正在开发Kafka Streams应用,因内部需求,将commit.interval.ms设置为极大值(如Long.MAX_VALUE)以禁用自动提交,转而通过底层Processor API在应用中显式触发提交操作。

运行过程中发现,应用使用的repartition topic中的记录始终未被删除,数据量持续无限增长。经过排查发现:与定期自动提交的行为不同,当手动触发提交时,StreamThread不会尝试清理repartition topic的记录。我查看了Kafka源码,该逻辑对应的位置在org.apache.kafka.streams.processor.internals.StreamThread.java的第1809行。

我想确认这是否是Kafka Streams设计中的预期行为。


预期结果

无论提交是由用户手动请求触发,还是基于commit.interval.ms配置的自动提交触发,当repartition.purge.interval.ms设定的时长过后,repartition topic的记录都应被清理至当前已提交的偏移量位置。

实际结果

仅当由commit.interval.ms触发自动提交时,repartition topic的记录才会被清理。


复现步骤(使用Spring Cloud Stream)

代码实现

SomeProcessor.java

@Bean
public Function<KStream<String, String>, KStream<String, Set<String>>> sampleStream() {
    return inputStream ->
        inputStream
            .map((key, value) -> KeyValue.pair(String.valueOf(key.hashCode()), value)) // 执行键值转换
            .repartition()
            .process(() -> new Processor<String, String, String, Set<String>>() {
                private ProcessorContext<String, Set<String>> context;
                private final Map<String, Set<String>> valueMap =  new HashMap<>();

                @Override
                public void init(ProcessorContext<String, Set<String>> context) {
                    this.context = context;
                    context.schedule(Duration.ofMinutes(1), PunctuationType.WALL_CLOCK_TIME, this::forward);
                }

                @Override
                public void process(Record<String, String> record) {
                    this.valueMap.putIfAbsent(record.key(), new HashSet<>());
                    this.valueMap.computeIfPresent(record.key(), (k, v) -> {
                        v.add(record.value());
                        return v;
                    });
                }

                @Override
                public void close() {
                    this.forward(context.currentSystemTimeMs());
                }

                private void forward(long timestamp) {
                    valueMap.forEach((key, value) -> context.forward(new Record<>(key, value, timestamp)));
                    context.commit();
                }
            });
}

application.yaml

spring:
  cloud:
    stream:
      kafka.streams:
        bindings:
          sampleStream-in-0:
            consumer:
              configuration:
                commit.interval.ms: 9223372036854775807 # 禁用自动提交
                repartition.purge.interval.ms: 300000 # 5分钟

内容的提问来源于stack exchange,提问作者Junhyun Kim

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:12:25