From c1534770c12b54db426fc7fdfe0ff5cda5edcb88 Mon Sep 17 00:00:00 2001 From: RongtongJin Date: Tue, 9 Sep 2025 16:36:29 +0800 Subject: [PATCH 1/3] feat: improve rocksdb compaction filter factory resource management - Refactor ConsumeQueueCompactionFilterFactory to use LongSupplier instead of MessageStore - Add proper resource cleanup in ConsumeQueueRocksDBStorage.preShutdown() - Update RocksDBOptionsFactory to accept external compaction filter factory - Optimize write buffer size for CQ rocksdb performance This change improves resource management and reduces memory leaks by: 1. Decoupling compaction filter from MessageStore reference 2. Ensuring proper cleanup of native resources 3. Making compaction filter factory lifecycle manageable --- .../ConsumeQueueCompactionFilterFactory.java | 10 ++-- .../rocksdb/ConsumeQueueRocksDBStorage.java | 18 +++++-- .../store/rocksdb/RocksDBOptionsFactory.java | 49 ++++++++++--------- 3 files changed, 46 insertions(+), 31 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueCompactionFilterFactory.java b/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueCompactionFilterFactory.java index aa796c4d39e..f19fb9e2036 100644 --- a/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueCompactionFilterFactory.java +++ b/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueCompactionFilterFactory.java @@ -16,20 +16,20 @@ */ package org.apache.rocketmq.store.rocksdb; +import java.util.function.LongSupplier; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.logging.org.slf4j.Logger; import org.apache.rocketmq.logging.org.slf4j.LoggerFactory; -import org.apache.rocketmq.store.MessageStore; import org.rocksdb.AbstractCompactionFilter; import org.rocksdb.AbstractCompactionFilterFactory; import org.rocksdb.RemoveConsumeQueueCompactionFilter; public class ConsumeQueueCompactionFilterFactory extends AbstractCompactionFilterFactory { private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.ROCKSDB_LOGGER_NAME); - private final MessageStore messageStore; + private final LongSupplier minPhyOffsetSupplier; - public ConsumeQueueCompactionFilterFactory(final MessageStore messageStore) { - this.messageStore = messageStore; + public ConsumeQueueCompactionFilterFactory(final LongSupplier minPhyOffsetSupplier) { + this.minPhyOffsetSupplier = minPhyOffsetSupplier; } @Override @@ -39,7 +39,7 @@ public String name() { @Override public RemoveConsumeQueueCompactionFilter createCompactionFilter(final AbstractCompactionFilter.Context context) { - long minPhyOffset = this.messageStore.getMinPhyOffset(); + long minPhyOffset = this.minPhyOffsetSupplier.getAsLong(); LOGGER.info("manualCompaction minPhyOffset: {}, isFull: {}, isManual: {}", minPhyOffset, context.isFullCompaction(), context.isManualCompaction()); return new RemoveConsumeQueueCompactionFilter(minPhyOffset); diff --git a/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueRocksDBStorage.java b/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueRocksDBStorage.java index b04aeab6bd6..4392283c67c 100644 --- a/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueRocksDBStorage.java +++ b/store/src/main/java/org/apache/rocketmq/store/rocksdb/ConsumeQueueRocksDBStorage.java @@ -38,6 +38,8 @@ public class ConsumeQueueRocksDBStorage extends AbstractRocksDBStorage { private final MessageStore messageStore; private volatile ColumnFamilyHandle offsetCFHandle; + private ConsumeQueueCompactionFilterFactory compactionFilterFactory; + public ConsumeQueueRocksDBStorage(final MessageStore messageStore, final String dbPath) { super(dbPath); this.messageStore = messageStore; @@ -65,7 +67,9 @@ protected boolean postLoad() { final List cfDescriptors = new ArrayList<>(); - ColumnFamilyOptions cqCfOptions = RocksDBOptionsFactory.createCQCFOptions(this.messageStore); + this.compactionFilterFactory = new ConsumeQueueCompactionFilterFactory(messageStore::getMinPhyOffset); + + ColumnFamilyOptions cqCfOptions = RocksDBOptionsFactory.createCQCFOptions(this.messageStore, this.compactionFilterFactory); this.cfOptions.add(cqCfOptions); cfDescriptors.add(new ColumnFamilyDescriptor(RocksDB.DEFAULT_COLUMN_FAMILY, cqCfOptions)); @@ -84,7 +88,14 @@ protected boolean postLoad() { @Override protected void preShutdown() { - this.offsetCFHandle.close(); + if (this.offsetCFHandle != null) { + this.offsetCFHandle.close(); + } + + if (this.compactionFilterFactory != null) { + this.compactionFilterFactory.close(); + } + } public byte[] getCQ(final byte[] keyBytes) throws RocksDBException { @@ -95,7 +106,8 @@ public byte[] getOffset(final byte[] keyBytes) throws RocksDBException { return get(this.offsetCFHandle, this.totalOrderReadOptions, keyBytes); } - public List multiGet(final List cfhList, final List keys) throws RocksDBException { + public List multiGet(final List cfhList, + final List keys) throws RocksDBException { return multiGet(this.totalOrderReadOptions, cfhList, keys); } diff --git a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java index 5687d6a222d..aeaa9c7024e 100644 --- a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java +++ b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java @@ -41,17 +41,18 @@ public class RocksDBOptionsFactory { - public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageStore) { + public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageStore, + ConsumeQueueCompactionFilterFactory consumeQueueCompactionFilterFactory) { BlockBasedTableConfig blockBasedTableConfig = new BlockBasedTableConfig(). - setFormatVersion(5). - setIndexType(IndexType.kBinarySearch). - setDataBlockIndexType(DataBlockIndexType.kDataBlockBinaryAndHash). - setDataBlockHashTableUtilRatio(0.75). - setBlockSize(32 * SizeUnit.KB). - setMetadataBlockSize(4 * SizeUnit.KB). - setFilterPolicy(new BloomFilter(16, false)). - setCacheIndexAndFilterBlocks(false). - setCacheIndexAndFilterBlocksWithHighPriority(true). + setFormatVersion(5). + setIndexType(IndexType.kBinarySearch). + setDataBlockIndexType(DataBlockIndexType.kDataBlockBinaryAndHash). + setDataBlockHashTableUtilRatio(0.75). + setBlockSize(32 * SizeUnit.KB). + setMetadataBlockSize(4 * SizeUnit.KB). + setFilterPolicy(new BloomFilter(16, false)). + setCacheIndexAndFilterBlocks(false). + setCacheIndexAndFilterBlocksWithHighPriority(true). setPinL0FilterAndIndexBlocksInCache(false). setPinTopLevelIndexAndFilter(true). setBlockCache(new LRUCache(1024 * SizeUnit.MB, 8, false)). @@ -72,8 +73,9 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt .getRocksdbCompressionType(); CompressionType bottomMostCompressionType = CompressionType.getCompressionType(bottomMostCompressionTypeOpt); CompressionType compressionType = CompressionType.getCompressionType(compressionTypeOpt); + return columnFamilyOptions.setMaxWriteBufferNumber(4). - setWriteBufferSize(128 * SizeUnit.MB). + setWriteBufferSize(32 * SizeUnit.MB). // speed up cq rocksdb open test setMinWriteBufferNumberToMerge(1). setTableFormatConfig(blockBasedTableConfig). setMemTableConfig(new SkipListMemTableConfig()). @@ -82,19 +84,20 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt setNumLevels(7). setCompactionPriority(CompactionPriority.MinOverlappingRatio). setCompactionStyle(CompactionStyle.UNIVERSAL). - setCompactionOptionsUniversal(compactionOption). - setMaxCompactionBytes(100 * SizeUnit.GB). - setSoftPendingCompactionBytesLimit(100 * SizeUnit.GB). - setHardPendingCompactionBytesLimit(256 * SizeUnit.GB). - setLevel0FileNumCompactionTrigger(2). - setLevel0SlowdownWritesTrigger(8). - setLevel0StopWritesTrigger(10). - setTargetFileSizeBase(256 * SizeUnit.MB). - setTargetFileSizeMultiplier(2). - setMergeOperator(new StringAppendOperator()). - setCompactionFilterFactory(new ConsumeQueueCompactionFilterFactory(messageStore)). - setReportBgIoStats(true). + setCompactionOptionsUniversal(compactionOption). + setMaxCompactionBytes(100 * SizeUnit.GB). + setSoftPendingCompactionBytesLimit(100 * SizeUnit.GB). + setHardPendingCompactionBytesLimit(256 * SizeUnit.GB). + setLevel0FileNumCompactionTrigger(2). + setLevel0SlowdownWritesTrigger(8). + setLevel0StopWritesTrigger(10). + setTargetFileSizeBase(256 * SizeUnit.MB). + setTargetFileSizeMultiplier(2). + setMergeOperator(new StringAppendOperator()). + setCompactionFilterFactory(consumeQueueCompactionFilterFactory). + setReportBgIoStats(true). setOptimizeFiltersForHits(true); + } public static ColumnFamilyOptions createOffsetCFOptions() { From cf5b3c7894617edc5577ec33741454d1f3bd90c6 Mon Sep 17 00:00:00 2001 From: RongtongJin Date: Tue, 9 Sep 2025 16:42:27 +0800 Subject: [PATCH 2/3] Polish the code --- .../store/rocksdb/RocksDBOptionsFactory.java | 46 +++++++++---------- 1 file changed, 22 insertions(+), 24 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java index aeaa9c7024e..03a41019a9d 100644 --- a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java +++ b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java @@ -44,15 +44,15 @@ public class RocksDBOptionsFactory { public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageStore, ConsumeQueueCompactionFilterFactory consumeQueueCompactionFilterFactory) { BlockBasedTableConfig blockBasedTableConfig = new BlockBasedTableConfig(). - setFormatVersion(5). - setIndexType(IndexType.kBinarySearch). - setDataBlockIndexType(DataBlockIndexType.kDataBlockBinaryAndHash). - setDataBlockHashTableUtilRatio(0.75). - setBlockSize(32 * SizeUnit.KB). - setMetadataBlockSize(4 * SizeUnit.KB). - setFilterPolicy(new BloomFilter(16, false)). - setCacheIndexAndFilterBlocks(false). - setCacheIndexAndFilterBlocksWithHighPriority(true). + setFormatVersion(5). + setIndexType(IndexType.kBinarySearch). + setDataBlockIndexType(DataBlockIndexType.kDataBlockBinaryAndHash). + setDataBlockHashTableUtilRatio(0.75). + setBlockSize(32 * SizeUnit.KB). + setMetadataBlockSize(4 * SizeUnit.KB). + setFilterPolicy(new BloomFilter(16, false)). + setCacheIndexAndFilterBlocks(false). + setCacheIndexAndFilterBlocksWithHighPriority(true). setPinL0FilterAndIndexBlocksInCache(false). setPinTopLevelIndexAndFilter(true). setBlockCache(new LRUCache(1024 * SizeUnit.MB, 8, false)). @@ -73,9 +73,8 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt .getRocksdbCompressionType(); CompressionType bottomMostCompressionType = CompressionType.getCompressionType(bottomMostCompressionTypeOpt); CompressionType compressionType = CompressionType.getCompressionType(compressionTypeOpt); - return columnFamilyOptions.setMaxWriteBufferNumber(4). - setWriteBufferSize(32 * SizeUnit.MB). // speed up cq rocksdb open test + setWriteBufferSize(128 * SizeUnit.MB). setMinWriteBufferNumberToMerge(1). setTableFormatConfig(blockBasedTableConfig). setMemTableConfig(new SkipListMemTableConfig()). @@ -84,20 +83,19 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt setNumLevels(7). setCompactionPriority(CompactionPriority.MinOverlappingRatio). setCompactionStyle(CompactionStyle.UNIVERSAL). - setCompactionOptionsUniversal(compactionOption). - setMaxCompactionBytes(100 * SizeUnit.GB). - setSoftPendingCompactionBytesLimit(100 * SizeUnit.GB). - setHardPendingCompactionBytesLimit(256 * SizeUnit.GB). - setLevel0FileNumCompactionTrigger(2). - setLevel0SlowdownWritesTrigger(8). - setLevel0StopWritesTrigger(10). - setTargetFileSizeBase(256 * SizeUnit.MB). - setTargetFileSizeMultiplier(2). - setMergeOperator(new StringAppendOperator()). - setCompactionFilterFactory(consumeQueueCompactionFilterFactory). - setReportBgIoStats(true). + setCompactionOptionsUniversal(compactionOption). + setMaxCompactionBytes(100 * SizeUnit.GB). + setSoftPendingCompactionBytesLimit(100 * SizeUnit.GB). + setHardPendingCompactionBytesLimit(256 * SizeUnit.GB). + setLevel0FileNumCompactionTrigger(2). + setLevel0SlowdownWritesTrigger(8). + setLevel0StopWritesTrigger(10). + setTargetFileSizeBase(256 * SizeUnit.MB). + setTargetFileSizeMultiplier(2). + setMergeOperator(new StringAppendOperator()). + setCompactionFilterFactory(new ConsumeQueueCompactionFilterFactory(messageStore::getMinPhyOffset)). + setReportBgIoStats(true). setOptimizeFiltersForHits(true); - } public static ColumnFamilyOptions createOffsetCFOptions() { From ec531c4b65c1249784d08804f4840ec2f6b9d324 Mon Sep 17 00:00:00 2001 From: RongtongJin Date: Wed, 10 Sep 2025 11:42:58 +0800 Subject: [PATCH 3/3] Fix consumeQueueCompactionFilterFactory not use --- .../apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java index 03a41019a9d..e365326c76d 100644 --- a/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java +++ b/store/src/main/java/org/apache/rocketmq/store/rocksdb/RocksDBOptionsFactory.java @@ -93,7 +93,7 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt setTargetFileSizeBase(256 * SizeUnit.MB). setTargetFileSizeMultiplier(2). setMergeOperator(new StringAppendOperator()). - setCompactionFilterFactory(new ConsumeQueueCompactionFilterFactory(messageStore::getMinPhyOffset)). + setCompactionFilterFactory(consumeQueueCompactionFilterFactory). setReportBgIoStats(true). setOptimizeFiltersForHits(true); }