[ISSUE #8976] Modify file segment construct method (#8977)

This commit is contained in:
lizhimins
2024-11-26 14:16:38 +08:00
committed by GitHub
parent a2eafb2ad5
commit fc2283008b
10 changed files with 80 additions and 54 deletions
@@ -17,7 +17,9 @@
package org.apache.rocketmq.tieredstore.file;
import com.google.common.annotations.VisibleForTesting;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.metadata.MetadataStore;
import org.apache.rocketmq.tieredstore.provider.FileSegmentFactory;
@@ -28,10 +30,19 @@ public class FlatFileFactory {
private final MessageStoreConfig storeConfig;
private final FileSegmentFactory fileSegmentFactory;
@VisibleForTesting
public FlatFileFactory(MetadataStore metadataStore, MessageStoreConfig storeConfig) {
this.metadataStore = metadataStore;
this.storeConfig = storeConfig;
this.fileSegmentFactory = new FileSegmentFactory(metadataStore, storeConfig);
this.fileSegmentFactory = new FileSegmentFactory(metadataStore, storeConfig, new MessageStoreExecutor());
}
public FlatFileFactory(MetadataStore metadataStore,
MessageStoreConfig storeConfig, MessageStoreExecutor executor) {
this.metadataStore = metadataStore;
this.storeConfig = storeConfig;
this.fileSegmentFactory = new FileSegmentFactory(metadataStore, storeConfig, executor);
}
public MessageStoreConfig getStoreConfig() {
@@ -50,7 +50,7 @@ public class FlatFileStore {
this.storeConfig = storeConfig;
this.metadataStore = metadataStore;
this.executor = executor;
this.flatFileFactory = new FlatFileFactory(metadataStore, storeConfig);
this.flatFileFactory = new FlatFileFactory(metadataStore, storeConfig, executor);
this.flatFileConcurrentMap = new ConcurrentHashMap<>();
}
@@ -23,6 +23,7 @@ import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Semaphore;
import java.util.concurrent.locks.ReentrantLock;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.AppendResult;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.exception.TieredStoreErrorCode;
@@ -45,6 +46,7 @@ public abstract class FileSegment implements Comparable<FileSegment>, FileSegmen
protected final MessageStoreConfig storeConfig;
protected final long maxSize;
protected final MessageStoreExecutor executor;
protected final ReentrantLock fileLock = new ReentrantLock();
protected final Semaphore commitLock = new Semaphore(1);
@@ -58,13 +60,14 @@ public abstract class FileSegment implements Comparable<FileSegment>, FileSegmen
protected volatile FileSegmentInputStream fileSegmentInputStream;
protected volatile CompletableFuture<Boolean> flightCommitRequest;
public FileSegment(MessageStoreConfig storeConfig,
FileSegmentType fileType, String filePath, long baseOffset) {
public FileSegment(MessageStoreConfig storeConfig, FileSegmentType fileType,
String filePath, long baseOffset, MessageStoreExecutor executor) {
this.storeConfig = storeConfig;
this.fileType = fileType;
this.filePath = filePath;
this.baseOffset = baseOffset;
this.executor = executor;
this.maxSize = this.getMaxSizeByFileType();
}
@@ -19,6 +19,7 @@ package org.apache.rocketmq.tieredstore.provider;
import java.lang.reflect.Constructor;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.metadata.MetadataStore;
@@ -26,16 +27,20 @@ public class FileSegmentFactory {
private final MetadataStore metadataStore;
private final MessageStoreConfig storeConfig;
private final MessageStoreExecutor executor;
private final Constructor<? extends FileSegment> fileSegmentConstructor;
public FileSegmentFactory(MetadataStore metadataStore, MessageStoreConfig storeConfig) {
public FileSegmentFactory(MetadataStore metadataStore,
MessageStoreConfig storeConfig, MessageStoreExecutor executor) {
try {
this.storeConfig = storeConfig;
this.metadataStore = metadataStore;
this.executor = executor;
Class<? extends FileSegment> clazz =
Class.forName(storeConfig.getTieredBackendServiceProvider()).asSubclass(FileSegment.class);
fileSegmentConstructor = clazz.getConstructor(
MessageStoreConfig.class, FileSegmentType.class, String.class, Long.TYPE);
MessageStoreConfig.class, FileSegmentType.class, String.class, Long.TYPE, MessageStoreExecutor.class);
} catch (Exception e) {
throw new RuntimeException(e);
}
@@ -51,7 +56,7 @@ public class FileSegmentFactory {
public FileSegment createSegment(FileSegmentType fileType, String filePath, long baseOffset) {
try {
return fileSegmentConstructor.newInstance(this.storeConfig, fileType, filePath, baseOffset);
return fileSegmentConstructor.newInstance(this.storeConfig, fileType, filePath, baseOffset, executor);
} catch (Exception e) {
throw new RuntimeException(e);
}
@@ -19,6 +19,7 @@ package org.apache.rocketmq.tieredstore.provider;
import java.nio.ByteBuffer;
import java.util.concurrent.CompletableFuture;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.stream.FileSegmentInputStream;
import org.apache.rocketmq.tieredstore.util.MessageStoreUtil;
@@ -35,9 +36,9 @@ public class MemoryFileSegment extends FileSegment {
protected boolean checkSize = true;
public MemoryFileSegment(MessageStoreConfig storeConfig,
FileSegmentType fileType, String filePath, long baseOffset) {
FileSegmentType fileType, String filePath, long baseOffset, MessageStoreExecutor executor) {
super(storeConfig, fileType, filePath, baseOffset);
super(storeConfig, fileType, filePath, baseOffset, executor);
memStore = ByteBuffer.allocate(10000);
memStore.position((int) getSize());
}
@@ -17,6 +17,7 @@
package org.apache.rocketmq.tieredstore.provider;
import com.google.common.base.Stopwatch;
import com.google.common.base.Supplier;
import com.google.common.io.ByteStreams;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.common.AttributesBuilder;
@@ -30,6 +31,7 @@ import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.metrics.TieredStoreMetricsManager;
import org.apache.rocketmq.tieredstore.stream.FileSegmentInputStream;
@@ -58,9 +60,9 @@ public class PosixFileSegment extends FileSegment {
private volatile FileChannel writeFileChannel;
public PosixFileSegment(MessageStoreConfig storeConfig,
FileSegmentType fileType, String filePath, long baseOffset) {
FileSegmentType fileType, String filePath, long baseOffset, MessageStoreExecutor executor) {
super(storeConfig, fileType, filePath, baseOffset);
super(storeConfig, fileType, filePath, baseOffset, executor);
// basePath
String basePath = StringUtils.defaultString(storeConfig.getTieredStoreFilePath(),
@@ -168,32 +170,30 @@ public class PosixFileSegment extends FileSegment {
AttributesBuilder attributesBuilder = newAttributesBuilder()
.put(LABEL_OPERATION, OPERATION_POSIX_READ);
CompletableFuture<ByteBuffer> future = new CompletableFuture<>();
ByteBuffer byteBuffer = ByteBuffer.allocate(length);
try {
readFileChannel.position(position);
readFileChannel.read(byteBuffer);
byteBuffer.flip();
byteBuffer.limit(length);
return CompletableFuture.supplyAsync((Supplier<ByteBuffer>) () -> {
ByteBuffer byteBuffer = ByteBuffer.allocate(length);
try {
readFileChannel.position(position);
readFileChannel.read(byteBuffer);
byteBuffer.flip();
byteBuffer.limit(length);
attributesBuilder.put(LABEL_SUCCESS, true);
long costTime = stopwatch.stop().elapsed(TimeUnit.MILLISECONDS);
TieredStoreMetricsManager.providerRpcLatency.record(costTime, attributesBuilder.build());
attributesBuilder.put(LABEL_SUCCESS, true);
long costTime = stopwatch.stop().elapsed(TimeUnit.MILLISECONDS);
TieredStoreMetricsManager.providerRpcLatency.record(costTime, attributesBuilder.build());
Attributes metricsAttributes = newAttributesBuilder()
.put(LABEL_OPERATION, OPERATION_POSIX_READ)
.build();
int downloadedBytes = byteBuffer.remaining();
TieredStoreMetricsManager.downloadBytes.record(downloadedBytes, metricsAttributes);
future.complete(byteBuffer);
} catch (IOException e) {
long costTime = stopwatch.stop().elapsed(TimeUnit.MILLISECONDS);
attributesBuilder.put(LABEL_SUCCESS, false);
TieredStoreMetricsManager.providerRpcLatency.record(costTime, attributesBuilder.build());
future.completeExceptionally(e);
}
return future;
Attributes metricsAttributes = newAttributesBuilder()
.put(LABEL_OPERATION, OPERATION_POSIX_READ)
.build();
int downloadedBytes = byteBuffer.remaining();
TieredStoreMetricsManager.downloadBytes.record(downloadedBytes, metricsAttributes);
} catch (IOException e) {
long costTime = stopwatch.stop().elapsed(TimeUnit.MILLISECONDS);
attributesBuilder.put(LABEL_SUCCESS, false);
TieredStoreMetricsManager.providerRpcLatency.record(costTime, attributesBuilder.build());
}
return byteBuffer;
}, executor.bufferFetchExecutor);
}
@Override
@@ -229,6 +229,6 @@ public class PosixFileSegment extends FileSegment {
return false;
}
return true;
});
}, executor.bufferCommitExecutor);
}
}
@@ -29,6 +29,7 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import org.apache.rocketmq.common.ThreadFactoryImpl;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.AppendResult;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.provider.FileSegment;
@@ -219,7 +220,7 @@ public class IndexStoreFileTest {
ByteBuffer byteBuffer = indexStoreFile.doCompaction();
FileSegment fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.INDEX, filePath, 0L);
storeConfig, FileSegmentType.INDEX, filePath, 0L, new MessageStoreExecutor());
fileSegment.append(byteBuffer, timestamp);
fileSegment.commitAsync().join();
Assert.assertEquals(byteBuffer.limit(), fileSegment.getSize());
@@ -252,7 +253,7 @@ public class IndexStoreFileTest {
ByteBuffer byteBuffer = indexStoreFile.doCompaction();
FileSegment fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.INDEX, filePath, 0L);
storeConfig, FileSegmentType.INDEX, filePath, 0L, new MessageStoreExecutor());
fileSegment.append(byteBuffer, timestamp);
fileSegment.commitAsync().join();
Assert.assertEquals(byteBuffer.limit(), fileSegment.getSize());
@@ -17,6 +17,7 @@
package org.apache.rocketmq.tieredstore.provider;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.metadata.DefaultMetadataStore;
import org.apache.rocketmq.tieredstore.metadata.MetadataStore;
@@ -34,9 +35,10 @@ public class FileSegmentFactoryTest {
MessageStoreConfig storeConfig = new MessageStoreConfig();
storeConfig.setTieredStoreCommitLogMaxSize(1024);
storeConfig.setTieredStoreFilePath(storePath);
MessageStoreExecutor executor = new MessageStoreExecutor();
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, executor);
Assert.assertEquals(metadataStore, factory.getMetadataStore());
Assert.assertEquals(storeConfig, factory.getStoreConfig());
@@ -60,7 +62,9 @@ public class FileSegmentFactoryTest {
() -> factory.createSegment(null, null, 0L));
storeConfig.setTieredBackendServiceProvider(null);
Assert.assertThrows(RuntimeException.class,
() -> new FileSegmentFactory(metadataStore, storeConfig));
() -> new FileSegmentFactory(metadataStore, storeConfig, executor));
executor.shutdown();
MessageStoreUtilTest.deleteStoreDirectory(storePath);
}
}
@@ -74,7 +74,7 @@ public class FileSegmentTest {
public void fileAttributesTest() {
int baseOffset = 1000;
FileSegment fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset, storeExecutor);
// for default value check
Assert.assertEquals(baseOffset, fileSegment.getBaseOffset());
@@ -104,9 +104,9 @@ public class FileSegmentTest {
@Test
public void fileSortByOffsetTest() {
FileSegment fileSegment1 = new PosixFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 200L);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 200L, storeExecutor);
FileSegment fileSegment2 = new PosixFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 100L);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 100L, storeExecutor);
FileSegment[] fileSegments = new FileSegment[] {fileSegment1, fileSegment2};
Arrays.sort(fileSegments);
Assert.assertEquals(fileSegments[0], fileSegment2);
@@ -116,17 +116,17 @@ public class FileSegmentTest {
@Test
public void fileMaxSizeTest() {
FileSegment fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 100L);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), 100L, storeExecutor);
Assert.assertEquals(storeConfig.getTieredStoreCommitLogMaxSize(), fileSegment.getMaxSize());
fileSegment.destroyFile();
fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.CONSUME_QUEUE, MessageStoreUtil.toFilePath(mq), 100L);
storeConfig, FileSegmentType.CONSUME_QUEUE, MessageStoreUtil.toFilePath(mq), 100L, storeExecutor);
Assert.assertEquals(storeConfig.getTieredStoreConsumeQueueMaxSize(), fileSegment.getMaxSize());
fileSegment.destroyFile();
fileSegment = new PosixFileSegment(
storeConfig, FileSegmentType.INDEX, MessageStoreUtil.toFilePath(mq), 100L);
storeConfig, FileSegmentType.INDEX, MessageStoreUtil.toFilePath(mq), 100L, storeExecutor);
Assert.assertEquals(Long.MAX_VALUE, fileSegment.getMaxSize());
fileSegment.destroyFile();
}
@@ -134,7 +134,7 @@ public class FileSegmentTest {
@Test
public void unexpectedCaseTest() {
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, new MessageStoreExecutor());
FileSegment fileSegment = factory.createCommitLogFileSegment(MessageStoreUtil.toFilePath(mq), baseOffset);
fileSegment.initPosition(fileSegment.getSize());
@@ -157,7 +157,7 @@ public class FileSegmentTest {
@Test
public void commitLogTest() throws InterruptedException {
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, new MessageStoreExecutor());
FileSegment fileSegment = factory.createCommitLogFileSegment(MessageStoreUtil.toFilePath(mq), baseOffset);
long lastSize = fileSegment.getSize();
fileSegment.initPosition(fileSegment.getSize());
@@ -225,7 +225,7 @@ public class FileSegmentTest {
@Test
public void consumeQueueTest() throws ClassNotFoundException, NoSuchMethodException {
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, new MessageStoreExecutor());
FileSegment fileSegment = factory.createConsumeQueueFileSegment(MessageStoreUtil.toFilePath(mq), baseOffset);
long storeTimestamp = System.currentTimeMillis();
@@ -258,7 +258,7 @@ public class FileSegmentTest {
@Test
public void fileSegmentReadTest() throws ClassNotFoundException, NoSuchMethodException {
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, new MessageStoreExecutor());
FileSegment fileSegment = factory.createConsumeQueueFileSegment(MessageStoreUtil.toFilePath(mq), baseOffset);
long storeTimestamp = System.currentTimeMillis();
@@ -292,7 +292,7 @@ public class FileSegmentTest {
@Test
public void commitFailedThenSuccessTest() {
MemoryFileSegment segment = new MemoryFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset, storeExecutor);
long lastSize = segment.getSize();
segment.setCheckSize(false);
@@ -352,7 +352,7 @@ public class FileSegmentTest {
public void commitFailedMoreTimes() {
long startTime = System.currentTimeMillis();
MemoryFileSegment segment = new MemoryFileSegment(
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset);
storeConfig, FileSegmentType.COMMIT_LOG, MessageStoreUtil.toFilePath(mq), baseOffset, storeExecutor);
long lastSize = segment.getSize();
segment.setCheckSize(false);
@@ -419,7 +419,7 @@ public class FileSegmentTest {
@Test
public void handleCommitExceptionTest() {
MetadataStore metadataStore = new DefaultMetadataStore(storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig);
FileSegmentFactory factory = new FileSegmentFactory(metadataStore, storeConfig, storeExecutor);
{
FileSegment fileSegment = factory.createCommitLogFileSegment(MessageStoreUtil.toFilePath(mq), baseOffset);
@@ -19,6 +19,7 @@ package org.apache.rocketmq.tieredstore.provider;
import java.io.IOException;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.tieredstore.MessageStoreConfig;
import org.apache.rocketmq.tieredstore.MessageStoreExecutor;
import org.apache.rocketmq.tieredstore.common.FileSegmentType;
import org.apache.rocketmq.tieredstore.stream.FileSegmentInputStream;
import org.apache.rocketmq.tieredstore.util.MessageStoreUtil;
@@ -34,7 +35,7 @@ public class MemoryFileSegmentTest {
public void memoryTest() throws IOException {
MemoryFileSegment fileSegment = new MemoryFileSegment(
new MessageStoreConfig(), FileSegmentType.COMMIT_LOG,
MessageStoreUtil.toFilePath(new MessageQueue()), 0L);
MessageStoreUtil.toFilePath(new MessageQueue()), 0L, new MessageStoreExecutor());
Assert.assertFalse(fileSegment.exists());
fileSegment.createFile();
MemoryFileSegment fileSpySegment = Mockito.spy(fileSegment);