This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git
The following commit(s) were added to refs/heads/develop by this push:
new e58f9fd35a [ISSUE #10959] Fix RocksDBConsumeQueue.iterateFrom should
reject offset below min offset (#10960)
e58f9fd35a is described below
commit e58f9fd35ad45b4bd11d749d2a368a3f6b3b9d38
Author: ymwneu <[email protected]>
AuthorDate: Wed Aug 19 10:56:51 2026 +0800
[ISSUE #10959] Fix RocksDBConsumeQueue.iterateFrom should reject offset
below min offset (#10960)
startIndex < minOffset was previously accepted as long as it stayed
within [0, maxOffset), returning a non-null LargeRocksDBConsumeQueueIterator
whose next() can yield a null CqUnit once the underlying data has been
purged. Callers like ScheduleMessageService rely on iterateFrom returning
null to detect and correct an out-of-range offset (matching ConsumeQueue's
getMinLogicOffset() check); without it they NPE dereferencing the null
CqUnit instead.
Add the same startIndex >= getMinOffsetInQueue() bound used by the
file-based ConsumeQueue to both iterateFrom overloads.
Co-authored-by: maowei.ymw <[email protected]>
---
.../rocketmq/store/queue/RocksDBConsumeQueue.java | 6 ++++--
.../store/queue/RocksDBConsumeQueueTest.java | 22 ++++++++++++++++++++++
2 files changed, 26 insertions(+), 2 deletions(-)
diff --git
a/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
b/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
index 86b4d3ef8b..19672d1b6b 100644
---
a/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
+++
b/store/src/main/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueue.java
@@ -298,8 +298,9 @@ public class RocksDBConsumeQueue implements
ConsumeQueueInterface {
@Override
public ReferredIterator<CqUnit> iterateFrom(final long startIndex) {
+ long minCqOffset = getMinOffsetInQueue();
long maxCqOffset = getMaxOffsetInQueue();
- if (startIndex < maxCqOffset && startIndex >= 0) {
+ if (startIndex >= minCqOffset && startIndex < maxCqOffset &&
startIndex >= 0) {
int num = pullNum(startIndex, maxCqOffset);
return new LargeRocksDBConsumeQueueIterator(startIndex, num);
}
@@ -308,8 +309,9 @@ public class RocksDBConsumeQueue implements
ConsumeQueueInterface {
@Override
public ReferredIterator<CqUnit> iterateFrom(long startIndex, int count)
throws RocksDBException {
+ long minCqOffset = getMinOffsetInQueue();
long maxCqOffset = getMaxOffsetInQueue();
- if (startIndex < maxCqOffset) {
+ if (startIndex >= minCqOffset && startIndex < maxCqOffset) {
int num = Math.min((int)(maxCqOffset - startIndex), count);
return iterateFrom0(startIndex, num, maxCqOffset);
}
diff --git
a/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
b/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
index 702d91fb07..f41d350437 100644
---
a/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
+++
b/store/src/test/java/org/apache/rocketmq/store/queue/RocksDBConsumeQueueTest.java
@@ -120,6 +120,28 @@ public class RocksDBConsumeQueueTest extends QueueTestBase
{
assertFalse(it.hasNext());
}
+ @Test
+ public void testIterateFrom_startIndexBelowMinOffset_returnsNull() throws
Exception {
+ if (MixAll.isMac()) {
+ return;
+ }
+ DefaultMessageStore messageStore = mock(DefaultMessageStore.class);
+ RocksDBConsumeQueueStore rocksDBConsumeQueueStore =
mock(RocksDBConsumeQueueStore.class);
+
when(messageStore.getQueueStore()).thenReturn(rocksDBConsumeQueueStore);
+ when(rocksDBConsumeQueueStore.getMinOffsetInQueue(anyString(),
anyInt())).thenReturn(9000L);
+ when(rocksDBConsumeQueueStore.getMaxOffsetInQueue(anyString(),
anyInt())).thenReturn(10000L);
+
+ RocksDBConsumeQueue consumeQueue = new
RocksDBConsumeQueue(messageStore.getMessageStoreConfig(),
rocksDBConsumeQueueStore, "topic", 0);
+
+ // startIndex has already been purged (below the queue's current min
offset); must not
+ // return a non-null iterator that later yields a null CqUnit and NPEs
its caller.
+ assertNull(consumeQueue.iterateFrom(8000));
+ assertNull(consumeQueue.iterateFrom(8000, 10));
+
+ // startIndex within [min, max) still works as before.
+ assertNotNull(consumeQueue.iterateFrom(9000));
+ }
+
@Test
public void testLmqCounter_running() throws ConsumeQueueException {
messageStore.getMessageStoreConfig().setEnableMultiDispatch(true);