From 134efd26a1265f18b805ef2dfa6f46a3ac3a0f56 Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:04:48 +0800 Subject: [PATCH 1/2] [ISSUE --- .../rocketmq/broker/pop/PopConsumerRocksdbStore.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java index dc68f9d9fe5..3e2fe8a36fe 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java @@ -147,9 +147,11 @@ public List scanExpiredRecords(long lower, long upper, int ma // configure prefix indexing to improve the performance of scans. // However, in the current implementation, this is not the bottleneck. List consumerRecordList = new ArrayList<>(); - try (ReadOptions scanOptions = new ReadOptions() - .setIterateLowerBound(new Slice(ByteBuffer.allocate(Long.BYTES).putLong(lower).array())) - .setIterateUpperBound(new Slice(ByteBuffer.allocate(Long.BYTES).putLong(upper).array())); + try (Slice lowerBound = new Slice(ByteBuffer.allocate(Long.BYTES).putLong(lower).array()); + Slice upperBound = new Slice(ByteBuffer.allocate(Long.BYTES).putLong(upper).array()); + ReadOptions scanOptions = new ReadOptions() + .setIterateLowerBound(lowerBound) + .setIterateUpperBound(upperBound); RocksIterator iterator = db.newIterator(this.columnFamilyHandle, scanOptions)) { iterator.seek(ByteBuffer.allocate(Long.BYTES).putLong(lower).array()); while (iterator.isValid() && consumerRecordList.size() < maxCount) { From 3da362b6cae0d89d5a4d38601b9b347bdb27d18c Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:05:52 +0800 Subject: [PATCH 2/2] [ISSUE #10898] Wake NettyEventExecutor immediately on shutdown --- .../rocketmq/remoting/netty/NettyRemotingAbstract.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java index a735f8455d3..b9bf58ccb44 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java @@ -784,7 +784,7 @@ public void run() { while (!this.isStopped()) { try { NettyEvent event = this.eventQueue.poll(3000, TimeUnit.MILLISECONDS); - if (event != null && listener != null) { + if (event != null && event.getType() != null && listener != null) { switch (event.getType()) { case IDLE: listener.onChannelIdle(event.getRemoteAddr(), event.getChannel()); @@ -814,6 +814,14 @@ public void run() { log.info(this.getServiceName() + " service end"); } + @Override + public void shutdown(final boolean interrupt) { + // wake up the thread blocked on eventQueue.poll() so it observes the stopped flag + // immediately instead of waiting for the next poll timeout + this.eventQueue.offer(new NettyEvent(null, null, null)); + super.shutdown(interrupt); + } + @Override public String getServiceName() { return NettyEventExecutor.class.getSimpleName();