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..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 @@ -41,7 +41,8 @@ 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). @@ -92,7 +93,7 @@ public static ColumnFamilyOptions createCQCFOptions(final MessageStore messageSt setTargetFileSizeBase(256 * SizeUnit.MB). setTargetFileSizeMultiplier(2). setMergeOperator(new StringAppendOperator()). - setCompactionFilterFactory(new ConsumeQueueCompactionFilterFactory(messageStore)). + setCompactionFilterFactory(consumeQueueCompactionFilterFactory). setReportBgIoStats(true). setOptimizeFiltersForHits(true); }