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 26cfb5f60e [ISSUE #10622] Resolve ConsumerOffsetManager dynamically in 
LiteEventDispatcher (#10623)
26cfb5f60e is described below

commit 26cfb5f60e47a943bd9d5d25b196c248222cc949
Author: Quan <[email protected]>
AuthorDate: Fri Jul 17 22:10:39 2026 +0800

    [ISSUE #10622] Resolve ConsumerOffsetManager dynamically in 
LiteEventDispatcher (#10623)
    
    - Remove cached consumerOffsetManager field from LiteEventDispatcher
    - Replace all usages with brokerController.getConsumerOffsetManager() for 
lazy resolution
    - Allows downstream projects to swap ConsumerOffsetManager after broker init
---
 .../java/org/apache/rocketmq/broker/lite/LiteEventDispatcher.java  | 7 ++-----
 1 file changed, 2 insertions(+), 5 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteEventDispatcher.java 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteEventDispatcher.java
index 2bd1f36186..ba6a62b38d 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteEventDispatcher.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteEventDispatcher.java
@@ -23,7 +23,6 @@ import com.google.common.cache.CacheBuilder;
 import org.apache.commons.collections.CollectionUtils;
 import org.apache.commons.lang3.tuple.Triple;
 import org.apache.rocketmq.broker.BrokerController;
-import org.apache.rocketmq.broker.offset.ConsumerOffsetManager;
 import org.apache.rocketmq.common.BrokerConfig;
 import org.apache.rocketmq.common.ServiceThread;
 import org.apache.rocketmq.common.constant.LoggerName;
@@ -63,7 +62,6 @@ public class LiteEventDispatcher extends ServiceThread {
     private final BrokerController brokerController;
     private final LiteSubscriptionRegistry liteSubscriptionRegistry;
     private final AbstractLiteLifecycleManager liteLifecycleManager;
-    private final ConsumerOffsetManager consumerOffsetManager;
 
     protected final ConcurrentMap<String, ClientEventSet> clientEventMap = new 
ConcurrentHashMap<>();
     protected final ConcurrentSkipListSet<FullDispatchRequest> fullDispatchSet 
= new ConcurrentSkipListSet<>(COMPARATOR);
@@ -78,7 +76,6 @@ public class LiteEventDispatcher extends ServiceThread {
         this.brokerController = brokerController;
         this.liteSubscriptionRegistry = liteSubscriptionRegistry;
         this.liteLifecycleManager = liteLifecycleManager;
-        this.consumerOffsetManager = 
brokerController.getConsumerOffsetManager();
     }
 
     public void init() {
@@ -230,7 +227,7 @@ public class LiteEventDispatcher extends ServiceThread {
             if (maxOffset <= 0) {
                 continue;
             }
-            long consumerOffset = consumerOffsetManager.queryOffset(group, 
lmqName, 0);
+            long consumerOffset = 
brokerController.getConsumerOffsetManager().queryOffset(group, lmqName, 0);
             if (consumerOffset >= maxOffset) {
                 continue;
             }
@@ -274,7 +271,7 @@ public class LiteEventDispatcher extends ServiceThread {
             if (maxOffset <= 0) {
                 return true;
             }
-            long consumerOffset = consumerOffsetManager.queryOffset(group, 
lmqName, 0);
+            long consumerOffset = 
brokerController.getConsumerOffsetManager().queryOffset(group, lmqName, 0);
             if (consumerOffset >= maxOffset) {
                 return true;
             }

Reply via email to