为何Kafka Streams消费速率始终是生产速率的两倍?
Kafka Streams消费速率为生产速率两倍的原因分析
问题描述
使用以下Prometheus指标监控Kafka Streams应用:
- 生产速率:
sum(irate(kafka_producer_producer_metrics_record_send_total{}[1m])) - 消费速率:
sum(irate(kafka_consumer_consumer_fetch_manager_metrics_records_consumed_total{}[1m]))
观测到所有Kafka Streams应用的消费速率始终是生产速率的两倍,既不存在生产速率反超消费速率的情况,内存也未出现暴涨现象。

补充配置信息
主题与业务场景
- 业务主题:1个分区,副本数为2
- Streams应用逻辑:仅包含简单的
map操作
Kafka集群配置
│ Kafka: │ Config: │ default.replication.factor: 3 │ inter.broker.protocol.version: 3.3 │ min.insync.replicas: 2 │ offsets.topic.replication.factor: 3 │ transaction.state.log.min.isr: 2 │ transaction.state.log.replication.factor: 3
Kafka Streams应用配置
props.put(StreamsConfig.TOPOLOGY_OPTIMIZATION_CONFIG, StreamsConfig.OPTIMIZE); props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1); props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 2); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest);
原因分析
这种消费速率恰好为生产速率两倍的现象,核心原因是Kafka Streams内部Changelog主题的消费被纳入了消费指标统计,具体逻辑如下:
Changelog主题的自动创建:即使你的应用只有简单的
map操作,但开启了TOPOLOGY_OPTIMIZATION_CONFIG=OPTIMIZE后,优化后的拓扑可能会自动引入状态存储(用于优化操作性能或隐含的状态依赖)。Kafka Streams会为每个状态存储创建对应的Changelog主题,用于备份状态数据,确保应用重启后能恢复状态。指标统计范围覆盖内部消费:你使用的消费指标
kafka_consumer_consumer_fetch_manager_metrics_records_consumed_total会统计所有消费者的拉取记录,包括Streams应用对业务输入主题和内部Changelog主题的消费。而生产指标仅统计业务主题的生产记录,两份消费数据叠加后,消费速率就刚好是生产速率的两倍。无内存暴涨的原因:Changelog主题的消费是为了同步状态存储的数据,Kafka Streams会自动处理状态的持久化与增量同步,不会无限制堆积数据,因此内存不会出现异常暴涨。
内容的提问来源于stack exchange,提问作者Konrad
相关产品推荐
相关产品推荐

