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

Kafka Streams应用内存占用异常过高问题排查求助

Kafka Streams应用内存占用远超预期排查问题

我在Kubernetes中运行一个Kafka Streams应用,仅将单个主题的数据读取到GlobalKTable中,再通过HTTP API提供键值对查询,但Pod内存占用远超预期,找不到原因。

我已根据Confluent官方文档配置了BoundedMemoryRocksDBConfigSetter,核心代码如下:

abstract class BoundedMemoryRocksDBConfig(blockCacheSize: Long = 50 * BoundedMemoryRocksDBConfig.Mi,
                                          numPartitions: Int = 10,
                                          ttl: Long = 180 * 24 * 60 * 60L)
  extends RocksDBConfigSetter
    with LazyLogging {
 
  private val totalOffHeapMemory = numPartitions * blockCacheSize
  private val totalMemTableMemory = totalOffHeapMemory
  assert(totalMemTableMemory <= totalOffHeapMemory)

  override def setConfig(storeName: String, options: Options, configs: util.Map[String, AnyRef]): Unit =
    BoundedMemoryRocksDBConfig.synchronized {
      if (!initialized) {
        sharedCache = Some(new org.rocksdb.LRUCache(totalOffHeapMemory))
        sharedWriteBufferManager = Some(new org.rocksdb.WriteBufferManager(totalOffHeapMemory, sharedCache.get))
        initialized = true
      }

      val tableConfig = options.tableFormatConfig.asInstanceOf[BlockBasedTableConfig]
      // These three options in combination will limit the memory used by RocksDB to the size passed
      // to the block cache (totalOffHeapMemory)
      tableConfig.setBlockCache(sharedCache.get)
      tableConfig.setCacheIndexAndFilterBlocks(true)
      options.setWriteBufferManager(sharedWriteBufferManager.get)
      options.setTableFormatConfig(tableConfig)

      val compactionsOptions = new CompactionOptionsFIFO()
      options.setCompactionOptionsFIFO(compactionsOptions)
      options.setTtl(ttl)
    }

  override def close(storeName: String, options: Options): Unit = {
    // do not close cache or writeBufferManager because they are shared among all stream tasks
  }
}
object BoundedMemoryRocksDBConfig {
  private var initialized = false
  private var sharedCache: Option[org.rocksdb.LRUCache] = None
  private var sharedWriteBufferManager: Option[org.rocksdb.WriteBufferManager] = None

  val Ki: Long = 1024L
  val Mi: Long = Ki * Ki
}

该主题包含10个分区,配置了50MB的LRUCache在块缓存和内存表/写缓冲区之间共享。此外,Kafka Streams默认的全拓扑记录缓存为10MB,还有生产者/消费者缓冲区、TCP收发缓冲区和可忽略大小的反序列化缓冲区。根据监控指标,Pod中的JVM提交了200MB堆内存,非堆内存使用量为160MB。按此计算,预期Pod内存占用最多为:10*50MB(块缓存+写缓冲区) + 10MB(记录缓存) + 360MB(JVM内存)= 870MB,但实际Pod内存占用远高于2GB且持续上升。

注意:Kafka Streams收集的以下RocksDB指标通常按task_id统计(此处为按task_id求和聚合)。我认为指标数值存在误差,实际应为当前数值的1/10——因为10个任务共享同一个LRUCache,但RocksDB不知道该缓存是共享的,因此每个任务单独报告500MB的块缓存使用量,而实际上所有任务共用的是同一块500MB的内存缓冲区。

Pod内存RSS
块缓存使用量
块缓存容量
所有内存表大小
预估表读取器内存
JVM非堆内存使用
JVM堆内存提交量

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 14:52:51