Kafka Streams 内存管理实战指南:Record Cache、RocksDB 与堆外内存调优
【免费下载链接】Kafka Apache Kafka – A distributed event streaming platform 项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
本文以 Apache Kafka Streams 官方开发者指南中的内存管理章节为骨架,系统讲解 Kafka Streams 应用中 record cache(记录缓存)、RocksDB 状态存储以及堆外(off-heap)内存的管理与调优方案。读完本文,你将掌握 cache.max.bytes.buffering / statestore.cache.max.bytes 与 commit.interval.ms 的协同机制、DSL 与 Processor API 中缓存行为的差异、RocksDB 内存占用与刷新(flush)频率对变更日志偏移量持久化和状态恢复性能的影响,以及通过 RocksDBConfigSetter 限定总堆外内存的完整做法。
说明:本仓库中 docs/documentation/streams/developer-guide/memory-mgmt.md 为 Hugo 站点重定向页,其实际内容位于 docs/streams/developer-guide/memory-mgmt.md,本文以其为正文基础并结合源码展开。
一、内存管理概览:缓存发生在写状态存储之前
Kafka Streams 允许你指定用于内部缓存和记录压缩(compacting)的总内存(RAM)大小。这段缓存发生在记录被写入状态存储(state store)或转发到下游节点之前,因此它直接决定了每个处理节点对外"可见"的记录数量,也决定了状态存储的写入压力。
需要特别注意的是,DSL 与 Processor API 中的 record cache 实现方式略有不同,后续章节将分别说明。此外,除了 record cache 本身,Kafka Streams 运行时的内存还来自 Kafka 客户端(producer/consumer 缓冲)、反序列化对象缓冲,以及最重要的 RocksDB 堆外内存,这些都在本文范围内逐一覆盖。
二、DSL 中的 Record Cache
2.1 哪些 KTable 会用到 record cache
在 DSL 中,你可以为整个处理拓扑(processing topology)实例指定 record cache 的总内存大小。它被以下两类 KTable 实例使用:
- Source KTable:通过 StreamsBuilder#table() 或 StreamsBuilder#globalTable() 创建的 KTable;
- Aggregation KTable:聚合(aggregations,参见 DSL API 聚合章节)产生的 KTable。
对于这些 KTable 实例,record cache 的作用有两个:
2.2 缓存与否的行为差异:一个聚合示例
假设输入是一个 KStream<String, Integer>,包含记录 <K,V>:<A, 1>, <D, 5>, <A, 20>, <A, 300>,我们只关注 key == A 的记录。一个聚合操作按 key 对 value 求和,输出 KTable<String, Integer>:
- 无缓存:对 key A 会输出一条表示聚合表变化的记录序列。括号 () 表示变化,左侧是新聚合值、右侧是旧聚合值:<A, (1, null)>, <A, (21, 1)>, <A, (321, 21)>——共 3 条记录;
- 有缓存:对 key A 只输出一条记录,中间记录在缓存中被压缩(compacted),最终输出 <A, (321, null)>,且只将该记录写入聚合内部状态存储并转发给下游操作。
这一对比直观说明:缓存带来的核心收益是显著减少写入状态存储和下游网络传输的记录量,从而降低磁盘 I/O 与 CPU 开销。
2.3 通过配置指定缓存大小
缓存大小通过 cache.max.bytes.buffering 参数指定,它是每个处理拓扑的全局设置:
// Enable record cache of size 10 MB.
Properties props = new Properties();
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
该参数控制分配给缓存的字节数。对于一个有 T 个线程、分配了 C 字节缓存的拓扑实例,每个线程将获得均等的 C/T 字节来构建自己的缓存,并可在其任务之间自由分配。这意味着缓存的个数等于线程数,线程之间不共享缓存。
2.4 缓存的内部机制:put / get 与 LRU 淘汰
缓存的基础 API 由 put() 和 get() 调用构成。当缓存达到大小上限后,记录按简单的 LRU(最近最少使用)方案被淘汰。处理节点处理完一条带 key 的记录 R1 = <K1, V1> 后,该记录在缓存中被标记为 dirty(脏);此后同一节点上处理到的任何同 key 记录 R2 = <K1, V2> 都会覆盖 <K1, V1>,这一行为称为"被压缩(being compacted)"。
这与 Kafka 的日志压缩(log compaction)效果相同,但发生得更早——记录还在客户端应用内存中时即被压缩,而非在服务端(Kafka broker)进行。flush 之后,R2 才会被转发到下一个处理节点并写入本地状态存储。
2.5 flush 时机:commit.interval.ms 与缓存压力的最早者优先
缓存的语义是:当 commit.interval.ms 与 cache.max.bytes.buffering(缓存压力)中最早者到达时,数据被 flush 到状态存储并转发到下游处理节点。这两个参数都是全局参数,因此无法为单个节点指定不同参数。
下面是根据目标场景的典型配置示例:
- 关闭缓存(缓存大小设为 0):
// Disable record cache
Properties props = new Properties();
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0);
- 开启缓存同时限制记录在缓存中的最长时间(设置 commit interval):
Properties props = new Properties();
// Enable record cache of size 10 MB.
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
// Set commit interval to 1 second.
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
下图为两个配置的联合效果示意(记录以蓝、红、黄、绿 4 个 key 展示,假设缓存只能容纳 3 个 key):

- 缓存禁用时(图 a):所有输入记录都会输出;
- 缓存启用时(图 b):
- 大部分记录在 commit interval 结束时输出(例如 t1 时刻输出一条蓝色记录,它是该时刻前 blue key 的最后一次覆盖写);
- 部分记录因缓存压力(即 commit interval 结束之前)提前输出,例如 t2 之前的红色记录。缓存越小,缓存压力越可能成为决定输出时机的首要因素;缓存越大,commit interval 则成为首要因素;
- 输出记录总数从 15 条减少到 8 条。
2.6 源码层面的配置证据
在 streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java 中可以确认这些配置的定义:
- CACHE_MAX_BYTES_BUFFERING_CONFIG = "cache.max.bytes.buffering"(第 472 行)——自 3.4 起已标记 @Deprecated,推荐改用 STATESTORE_CACHE_MAX_BYTES_CONFIG = "statestore.cache.max.bytes"(第 810 行);
- BUFFERED_RECORDS_PER_PARTITION_CONFIG = "buffered.records.per.partition"(第 460 行)——控制反序列化对象缓冲,后文会再次提到;
- COMMIT_INTERVAL_MS_CONFIG = "commit.interval.ms"(第 484 行)——提交处理进度的频率;在 exactly-once(processing.guarantee 为 EXACTLY_ONCE_V2)模式下默认值不同,源码注释指出其默认值为 EOS_DEFAULT_COMMIT_INTERVAL_MS,否则为 DEFAULT_COMMIT_INTERVAL_MS;
- 从 3.4 起新增 STATESTORE_UNCOMMITTED_MAX_BYTES_CONFIG = "statestore.uncommitted.max.bytes"(第 815 行),默认 64 MB(67_108_864L),用于在存在事务性状态存储时,即使未到 commit.interval.ms,也通过"未提交字节数"上限触发提前提交,且限额按 stream thread 数均分(含 global state thread)。
三、Processor API 中的 Record Cache
在 Processor API 中,同样可以为拓扑实例指定 record cache 的总内存大小,但它仅用于在有状态 processor node 写入其状态存储之前对输出记录进行内部缓存和压缩。
与 DSL 的关键差异是:Processor API 的 record cache 不会缓存或压缩转发到下游的输出记录。这意味着:
- 所有下游 processor node 都能看到全部记录;
- 而状态存储只看到减少后的记录数量。
这不影响系统正确性,而是针对状态存储的一种性能优化。例如,使用 Processor API 时,你可以在把一条记录存入状态存储的同时,向下游转发一个不同的值。
沿用 Processor API 状态存储章节中的示例,可以通过 withCachingDisabled 调用关闭缓存(注意缓存默认启用,但 API 也显式提供了 withCachingEnabled 调用):
StoreBuilder countStoreBuilder =
Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("Counts"),
Serdes.String(),
Serdes.Long())
.withCachingEnabled();
需要特别指出的是,record cache 不支持 versioned state stores(版本化状态存储,参见 Processor API 状态存储-版本化章节)。
为避免读取到陈旧数据,可以在创建迭代器之前对 store 调用 flush()。但需要注意:过于频繁地 flush 在使用 RocksDB 时可能导致性能下降,因此通常建议避免手动 flush。
四、RocksDB:状态存储的堆外内存与刷新策略
Kafka Streams 的持久化状态存储基于 RocksDB。每个 RocksDB 实例都会分配堆外内存,用于 block cache(块缓存)、index 与 filter blocks(索引和过滤器块)以及 memtable(写缓冲区)。关键配置(以 RocksDB 4.1.0 版本为参照)包括 block_cache_size、write_buffer_size 和 max_write_buffer_number,这些都可以通过 rocksdb.config.setter 配置指定。
在源码 streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java 中可以看到默认值(第 106-112 行):WRITE_BUFFER_SIZE = 16MB、BLOCK_CACHE_SIZE = 50MB、BLOCK_SIZE = 4096B、MAX_WRITE_BUFFERS = 3,且默认 NO_COMPRESSION 压缩、UNIVERSAL 压缩风格。自定义设置通过 RocksDBConfigSetter 接口注入(setConfig / close 两个方法),其 Javadoc 特别提醒:若要修改 BlockBasedTableConfig,应通过 (BlockBasedTableConfig) options.tableFormatConfig() 获取现有实例再修改,以免丢失其他默认设置。
4.1 变更日志偏移量持久化与刷新频率
(4.3 / KIP-1035 起) Kafka Streams 将每个持久化 store 的变更日志(changelog)偏移量存储在 RocksDB 内部,而非原来每个任务一个的 .checkpoint 文件。同时 Kafka Streams 以**禁用 WAL(write-ahead log)**的方式运行 RocksDB——变更日志主题本身就是记录的持久化日志——因此数据和变更日志偏移量只有在 memtable flush 到 SST 文件时才真正落盘。flush 发生在以下两种时机:
- memtable 写满 write_buffer_size(默认 16 MB);
- store 干净关闭(KafkaStreams#close)。
与早期版本不同,Kafka Streams 不再在每次 commit 时强制 flush RocksDB。
实际影响:磁盘上的变更日志偏移量只和最近一次自然 flush 或干净关闭一样"新鲜"。
- 对高吞吐的 store,memtable 频繁写满,这通常不是问题;
- 对低流量 store,memtable 可能需要很长时间才能写满,持久化的偏移量会长期滞后于 store 的实际位置。如果应用随后非正常退出(例如 SIGKILL、OOM-kill,或 KafkaStreams#close 未能在进程/Pod 关闭宽限期内完成),只有最后一次 flush 的偏移量得以保留。若此时变更日志主题的 log-start offset 已经因保留策略或压缩推进到该陈旧偏移量之后,则下次重启时恢复(restore)消费者 seek 越界,任务将从变更日志重新初始化(日志中出现 OffsetOutOfRangeException / TaskCorruptedException)。这会自动恢复且无数据丢失,但会导致任务全量重新恢复(re-restore)。
缓解措施:针对受影响的低流量 store,可以通过自定义 RocksDBConfigSetter 降低其 write_buffer_size(并相应调整 max_write_buffer_number),使其更频繁地 flush:
public static class CustomRocksDBConfig implements RocksDBConfigSetter {
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
// smaller write buffer => more frequent flushes => fresher on-disk changelog offset,
// at the cost of more (smaller) SST files and more compaction work
options.setWriteBufferSize(4 * 1024 * 1024L);
}
@Override
public void close(final String storeName, final Options options) {}
}
这是一个权衡:更小的写缓冲区限制了持久化偏移量可能滞后的程度,但会增加 SST 文件数量和压缩负载。因此应只调优确实需要的低流量 store,并根据 store 的写入速率来设定缓冲区大小。
注意,flush 是基于写入量而非时间的:一个只有零星写入的 store 可能直到关闭时都不 flush,因此最可靠的保护是干净关闭。建议注册一个调用 KafkaStreams#close 的 shutdown hook,并确保它能完整执行——在 Kubernetes 上,应将 close 超时设置在 Pod 终止宽限期(默认 30s)之下留出余量,避免进程在 close 中途被 SIGKILL。
另外,store 前面的 record cache(statestore.cache.max.bytes)会把同一 key 的重复更新在原地合并后再写入 RocksDB。因此对于更新频繁、keyspace 小的工作负载,每个 commit 到达 memtable 的字节更少,memtable 填充得更慢。
4.2 恢复性能与写缓冲区大小
flush 同样影响状态恢复(restore)——恢复过程是批量写入 store 的。memtable 写满即产生一次 flush,每次 flush 产生一个 level-0(L0)SST 文件。因此,用较小的写缓冲区恢复大状态量的 store,会产生大量 L0 文件和大规模压缩工作,恢复可能因 L0 文件数(level0_slowdown_writes_trigger / level0_stop_writes_trigger)或 flush 队列(max_write_buffer_number)而停滞,导致恢复耗时更长、时长更不可预测。
如果恢复时间是需要关注的点,建议通过自定义 RocksDBConfigSetter 将 write_buffer_size 提高到 16 MB 默认值之上,起步范围 32 MB 到 64 MB 是合理的,例如:
options.setWriteBufferSize(64 * 1024 * 1024L);
更少但更大的 flush 意味着更少的 L0 文件和更少的压缩工作,通常恢复更快。进一步调高 level0_file_num_compaction_trigger 也可能有帮助。两个设置都会消耗堆外内存:一个 store 最多持有 max_write_buffer_number 个缓冲区,再乘以实例上的 store 数量与 stream thread 数。如果使用共享 WriteBufferManager 限定总内存(见下文),需要为此留出预算。
请注意:这条建议与上一节的建议方向相反,因为适用的场景不同——低流量 store(memtable 很少写满)应缩小缓冲区;恢复慢的 store 应扩大缓冲区。
4.3 堆外内存使用:jemalloc 与总量限制
推荐更换 RocksDB 默认内存分配器,因为默认分配器可能导致内存消耗增加。要切换到 jemalloc,需要在启动 Kafka Streams 应用前设置环境变量 LD_PRELOAD:
# example: install jemalloc (on Debian)
$ apt install -y libjemalloc-dev
# set LD_PRELOAD before you start your Kafka Streams application
$ export LD_PRELOAD="/usr/lib/x86_64-linux-gnu/libjemalloc.so"
从 2.3.0 起,可以限定所有实例的总内存,从而限制 Kafka Streams 应用的堆外内存总量。做法是:配置 RocksDB 将 index 和 filter blocks 缓存到 block cache 中、通过共享的 WriteBufferManager 限制 memtable 内存并让其计入 block cache,然后将同一个 Cache 对象传给每个实例。一个实现此方案的示例 RocksDBConfigSetter 如下:
public static class BoundedMemoryRocksDBConfig implements RocksDBConfigSetter {
private static org.rocksdb.Cache cache = new org.rocksdb.LRUCache(TOTAL_OFF_HEAP_MEMORY, -1, false, INDEX_FILTER_BLOCK_RATIO); // (1)
private static org.rocksdb.WriteBufferManager writeBufferManager = new org.rocksdb.WriteBufferManager(TOTAL_MEMTABLE_MEMORY, cache);
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
BlockBasedTableConfig tableConfig = (BlockBasedTableConfig) options.tableFormatConfig();
// These three options in combination will limit the memory used by RocksDB to the size passed to the block cache (TOTAL_OFF_HEAP_MEMORY)
tableConfig.setBlockCache(cache);
tableConfig.setCacheIndexAndFilterBlocks(true);
options.setWriteBufferManager(writeBufferManager);
// These options are recommended to be set when bounding the total memory
tableConfig.setCacheIndexAndFilterBlocksWithHighPriority(true); // (2)
tableConfig.setPinTopLevelIndexAndFilter(true);
tableConfig.setBlockSize(BLOCK_SIZE); // (3)
options.setMaxWriteBufferNumber(N_MEMTABLES);
options.setWriteBufferSize(MEMTABLE_SIZE);
options.setTableFormatConfig(tableConfig);
}
@Override
public void close(final String storeName, final Options options) {
// Cache and WriteBufferManager should not be closed here, as the same objects are shared by every store instance.
}
}
脚注说明:
注意:虽然上述配置是推荐的最低集合,但具体哪种组合性能最佳取决于工作负载,建议结合自己的用例进行实验。一个应用的最优配置未必适用于拓扑或输入主题不同的另一个应用。除了上述推荐配置,还可以考虑 RocksDB 的分区索引过滤器(partitioned index filters)。
五、其他内存使用:客户端缓冲与反序列化对象
Apache Kafka 中还有其它模块在运行时分配内存,包括:
- Producer 缓冲:由 producer 配置 buffer.memory 管理;
- Consumer 缓冲:目前没有严格管理,但可以通过 fetch 大小间接控制,即 fetch.max.bytes 和 fetch.max.wait.ms;
- TCP 收发缓冲:producer 和 consumer 各自有独立的 TCP send / receive 缓冲区,不计入上述 buffering 内存,由 send.buffer.bytes / receive.buffer.bytes 配置控制;
- 反序列化对象缓冲:consumer.poll() 返回记录后,记录会被反序列化以提取时间戳并缓存在 streams 空间中。目前只能通过 buffered.records.per.partition 间接控制(对应源码 StreamsConfig.java 第 460 行的 BUFFERED_RECORDS_PER_PARTITION_CONFIG)。
六、易被忽视的资源释放:迭代器必须显式关闭
迭代器应显式关闭以释放资源:store 迭代器(例如 KeyValueIterator 和 WindowStoreIterator)必须在用完后显式关闭,以释放打开的文件句柄和内存读缓冲区等资源,或者对这类 Closeable 类使用 try-with-resources 语句(JDK 7 起可用)。
否则,stream 应用的内存使用会随着运行持续增长,直至触发 OOM。这一点在长时间运行的生产应用中尤为关键,建议将迭代器访问统一收敛到 try-with-resources 模式。
七、总结与调优路径
Kafka Streams 的内存管理可以归纳为三个层次,对应不同的配置手段:
| Record cache(记录缓存) | cache.max.bytes.buffering(3.4 起建议 statestore.cache.max.bytes)、commit.interval.ms | 在写入状态存储与转发下游前压缩重复 key 更新;flush 时机由两者最早者决定;按线程均分缓存(C/T) |
| RocksDB 堆外内存 | rocksdb.config.setter + RocksDBConfigSetter:write_buffer_size、max_write_buffer_number、block cache、WriteBufferManager | 控制 memtable/flush 频率(影响 changelog 偏移量新鲜度与恢复性能)、限制总堆外内存、LD_PRELOAD 启用 jemalloc |
| 客户端与反序列化缓冲 | buffer.memory、fetch.max.bytes、fetch.max.wait.ms、send.buffer.bytes、receive.buffer.bytes、buffered.records.per.partition | 管理 producer/consumer 缓冲、TCP 收发缓冲与反序列化对象缓冲 |
实践要点:
更多相关内容可继续阅读 Kafka Streams 开发者指南、DSL API 文档 与 Processor API 文档。
【免费下载链接】Kafka Apache Kafka – A distributed event streaming platform 项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考






