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
相关产品推荐
相关产品推荐

