diff --git a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/file/FlatAppendFile.java b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/file/FlatAppendFile.java index 6380d90a491..c8ccc410efd 100644 --- a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/file/FlatAppendFile.java +++ b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/file/FlatAppendFile.java @@ -241,23 +241,30 @@ public CompletableFuture readAsync(long offset, int length) { } } - FileSegment fileSegment1 = fileSegmentList.get(index); - FileSegment fileSegment2 = offset + length > fileSegment1.getCommitOffset() && - fileSegmentList.size() > index + 1 ? fileSegmentList.get(index + 1) : null; + FileSegment fileSegment = fileSegmentList.get(index); + if (offset + length <= fileSegment.getCommitOffset() || fileSegmentList.size() <= index + 1) { + return fileSegment.readAsync(offset - fileSegment.getBaseOffset(), length); + } - if (fileSegment2 == null) { - return fileSegment1.readAsync(offset - fileSegment1.getBaseOffset(), length); + List> futureList = new ArrayList<>(); + long readOffset = offset; + int remainingLength = length; + for (; index < fileSegmentList.size() && remainingLength > 0; index++) { + fileSegment = fileSegmentList.get(index); + int segmentLength = (int) Math.min(remainingLength, fileSegment.getCommitOffset() - readOffset); + futureList.add(fileSegment.readAsync(readOffset - fileSegment.getBaseOffset(), segmentLength)); + readOffset += segmentLength; + remainingLength -= segmentLength; } - int segment1Length = (int) (fileSegment1.getCommitOffset() - offset); - return fileSegment1.readAsync(offset - fileSegment1.getBaseOffset(), segment1Length) - .thenCombine(fileSegment2.readAsync(0, length - segment1Length), - (buffer1, buffer2) -> { - ByteBuffer buffer = ByteBuffer.allocate(buffer1.remaining() + buffer2.remaining()); - buffer.put(buffer1).put(buffer2); - buffer.flip(); - return buffer; - }); + CompletableFuture[] futures = futureList.toArray(new CompletableFuture[0]); + return CompletableFuture.allOf(futures).thenApply(nil -> { + int resultLength = futureList.stream().mapToInt(future -> future.join().remaining()).sum(); + ByteBuffer result = ByteBuffer.allocate(resultLength); + futureList.forEach(future -> result.put(future.join())); + result.flip(); + return result; + }); } public void shutdown() { diff --git a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/file/FlatAppendFileTest.java b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/file/FlatAppendFileTest.java index b3df4e8aece..e2859b42e03 100644 --- a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/file/FlatAppendFileTest.java +++ b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/file/FlatAppendFileTest.java @@ -198,6 +198,30 @@ public void testAppendAndRead() { flatFile.destroy(); } + @Test + public void testReadAcrossMultipleFileSegments() { + storeConfig.setTieredStoreConsumeQueueMaxSize(100L); + FlatAppendFile flatFile = flatFileFactory.createFlatFileForConsumeQueue(MessageStoreUtil.toFilePath(queue)); + flatFile.rollingNewFile(0L); + + for (byte value = 1; value <= 3; value++) { + byte[] data = new byte[100]; + Arrays.fill(data, value); + flatFile.append(ByteBuffer.wrap(data), 1L); + } + flatFile.commitAsync().join(); + + ByteBuffer result = flatFile.readAsync(50L, 250).join(); + byte[] actual = new byte[result.remaining()]; + result.get(actual); + byte[] expected = new byte[250]; + Arrays.fill(expected, 0, 50, (byte) 1); + Arrays.fill(expected, 50, 150, (byte) 2); + Arrays.fill(expected, 150, 250, (byte) 3); + Assert.assertArrayEquals(expected, actual); + flatFile.destroy(); + } + @Test public void testCleanExpiredFile() { FlatAppendFile flatFile = flatFileFactory.createFlatFileForConsumeQueue(MessageStoreUtil.toFilePath(queue));