From 94fbfcf342fe13fce9cba4a6c31d052aaaf920e1 Mon Sep 17 00:00:00 2001 From: lizhimins <707364882@qq.com> Date: Tue, 21 Apr 2026 19:15:26 +0800 Subject: [PATCH] [ISSUE #10268] Fix incorrect time range file selection in IndexStoreService.queryAsync (#10269) Co-authored-by: lizhimins --- .../tieredstore/index/IndexStoreService.java | 5 ++- .../index/IndexStoreServiceTest.java | 37 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreService.java b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreService.java index 132d2162f9..bf91e051ea 100644 --- a/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreService.java +++ b/tieredstore/src/main/java/org/apache/rocketmq/tieredstore/index/IndexStoreService.java @@ -235,11 +235,14 @@ public class IndexStoreService extends ServiceThread implements IndexService { try { readWriteLock.readLock().lock(); ConcurrentNavigableMap pendingMap = - this.timeStoreTable.subMap(beginTime, true, endTime, true); + this.timeStoreTable.headMap(endTime, true); List> futureList = new ArrayList<>(pendingMap.size()); ConcurrentSkipListMap result = new ConcurrentSkipListMap<>(); for (Map.Entry entry : pendingMap.descendingMap().entrySet()) { + if (entry.getValue().getEndTimestamp() < beginTime) { + break; + } CompletableFuture completableFuture = entry.getValue() .queryAsync(topic, key, maxCount, beginTime, endTime) .thenAccept(itemList -> itemList.forEach(indexItem -> { diff --git a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/index/IndexStoreServiceTest.java b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/index/IndexStoreServiceTest.java index 7b881ddd44..90b706a96f 100644 --- a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/index/IndexStoreServiceTest.java +++ b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/index/IndexStoreServiceTest.java @@ -351,4 +351,41 @@ public class IndexStoreServiceTest { executorService.shutdown(); Assert.assertTrue(result.get()); } + + @Test + public void queryCrossFileBoundaryTest() throws InterruptedException, ExecutionException { + indexService = new IndexStoreService(fileAllocator, filePath); + indexService.start(); + + // Create first file with early beginTime + long file1Begin = System.currentTimeMillis(); + for (int i = 0; i < storeConfig.getTieredStoreIndexFileMaxIndexNum() - 1; i++) { + indexService.putKey(TOPIC_NAME, TOPIC_ID, QUEUE_ID, + Collections.singleton("crossKey"), i * 100L, MESSAGE_SIZE, file1Begin + i * 1000); + } + + // Create second file with later beginTime (beyond query range) + long file2Begin = System.currentTimeMillis() + 100_000; + indexService.createNewIndexFile(file2Begin); + for (int i = 0; i < 5; i++) { + indexService.putKey(TOPIC_NAME, TOPIC_ID, QUEUE_ID, + Collections.singleton("crossKey"), i * 100L, MESSAGE_SIZE, file2Begin + i); + } + + Assert.assertEquals(2, indexService.getTimeStoreTable().size()); + + // Query range: beginTime is AFTER file1's beginTime but BEFORE file1's last item timestamp + // This should select file1, NOT file2 (file2 beginTime > queryEnd) + long queryBegin = file1Begin + 5_000; + long queryEnd = file1Begin + 15_000; + + List results = indexService.queryAsync( + TOPIC_NAME, "crossKey", 10, queryBegin, queryEnd).get(); + + // file1 has items at timestamps: file1Begin, file1Begin+1000, ..., file1Begin+(N-1)*1000 + // Items in range [file1Begin+5000, file1Begin+15000] should match + // The bug (subMap) would return empty because file1's key < queryBegin + Assert.assertFalse("Should find index items from file covering query range", results.isEmpty()); + Assert.assertTrue("Should find items within query time range", results.size() > 0); + } } \ No newline at end of file