I am trying to make a simple Kafka Stream application (v2.3.1 initially it was 2.3.0) which gathers statistics on specified time intervals (e.g. per minute namely tumbling windows). Therefore I follow a textbook implementation like the one below
events
.groupByKey()
.windowedBy(
TimeWindows.of(Duration.ofMinutes(1).grace(Duration.ZERO))
aggregate(...),Materialized.as("agg-metric")).withRetention(Duration.ofMinutes(5))
.suppress(Suppressed.untilWindowClose(BufferConfig.unbounded()))
.toStream((key, value) -> key.key())
Everything seems to works normally except that my memory footprint is constantly growing. I can see that my heap memory is stable so I assume that the issue is related with RocksDB instances created.
I have 10 partitions and since each window (with a default of 3 segments per partition), creates 3 segments, I expect in total to have 30 instances of RocksDB. Since the default configuration values for RocksDB are rather large for my app I choose to change the default configuration for RocksDB, and implement RocksDBConfigSetter according to code below essentially trying to impose a limit on off-heap memory consumption.
private static final long BLOCK_CACHE_SIZE = 8 * 1024 * 1024L;
private static final long BLOCK_SIZE = 4096L;
private static final long WRITE_BUFFER_SIZE = 2 * 1024 * 1024L;
private static final int MAX_WRITE_BUFFERS = 2;
private org.rocksdb.Cache cache = new org.rocksdb.LRUCache(BLOCK_CACHE_SIZE); // 8MB
private org.rocksdb.Filter filter = new org.rocksdb.BloomFilter();
@Override
public void setConfig(final String storeName, final Options options, final Map<String, Object> configs) {
BlockBasedTableConfig tableConfig = new org.rocksdb.BlockBasedTableConfig();
tableConfig.setBlockCache(cache) // 8 MB
tableConfig.setBlockSize(BLOCK_SIZE); // 4 KB
tableConfig.setCacheIndexAndFilterBlocks(true);;
tableConfig.setPinTopLevelIndexAndFilter(true);
tableConfig.setFilter(new org.rocksdb.BloomFilter())
options.setMaxWriteBufferNumber(MAX_WRITE_BUFFERS); // 2 memtables
options.setWriteBufferSize(WRITE_BUFFER_SIZE); // 2 MB
options.setInfoLogLevel(InfoLogLevel.INFO_LEVEL);
options.setTableFormatConfig(tableConfig);
}
Based on the configuration values above and https://docs.confluent.io/current/streams/sizing.html I expect that the total memory allocated to RocksDB would be
(write_buffer_size_mb * write_buffer_count) + block_cache_size_mb => 2*2 + 8 => 12MB
Therefore for 30 instances the total allocated off-heap memory would account to 12 * 30 = 360 MB.
When I try to run this app on a VM with 2G memory I assign 512MB for the heap of the kafka streams app, so based on my logic/understanding the total memory allocated should plateau at a value lower than 1 GB (512 + 360).
Unfortunately this does not seem to be the case since my memory does not stop to grow, albeit slowly after a point, but steadily almost ~2% per day and it is unavoidable to consume all VM memory at some point and finally kill the process. The more concerning fact is that I never witness any release of the off-heap memory, even when my traffic is getting very low.
As a result, I am wondering what am I doing wrong in such a simple and common use-case. Am I missing something while calculating memory consumption for my app ? Is it possible to limit my memory consumption on my VM and what settings do I need to change on my configuration to limit memory allocation for Kafka stream app and RocksDB instances ?