Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<RemoveConsumeQueueCompactionFilter> {
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
Expand All @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -65,7 +67,9 @@ protected boolean postLoad() {

final List<ColumnFamilyDescriptor> 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));

Expand All @@ -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 {
Expand All @@ -95,7 +106,8 @@ public byte[] getOffset(final byte[] keyBytes) throws RocksDBException {
return get(this.offsetCFHandle, this.totalOrderReadOptions, keyBytes);
}

public List<byte[]> multiGet(final List<ColumnFamilyHandle> cfhList, final List<byte[]> keys) throws RocksDBException {
public List<byte[]> multiGet(final List<ColumnFamilyHandle> cfhList,
final List<byte[]> keys) throws RocksDBException {
return multiGet(this.totalOrderReadOptions, cfhList, keys);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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);
}
Expand Down
Loading