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;
}