diff --git a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java index 3d4994f533..454b96cd5e 100644 --- a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java +++ b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java @@ -401,10 +401,6 @@ public class BrokerConfig extends BrokerIdentity { */ private boolean estimateAccumulation = true; - /** - * Build ConsumeQueue concurrently with multi-thread - */ - private boolean enableBuildConsumeQueueConcurrently = false; public long getMaxPopPollingSize() { return maxPopPollingSize; @@ -1661,12 +1657,4 @@ public class BrokerConfig extends BrokerIdentity { public void setBatchDispatchRequestThreadPoolNums(int batchDispatchRequestThreadPoolNums) { this.batchDispatchRequestThreadPoolNums = batchDispatchRequestThreadPoolNums; } - - public boolean isEnableBuildConsumeQueueConcurrently() { - return enableBuildConsumeQueueConcurrently; - } - - public void setEnableBuildConsumeQueueConcurrently(boolean enableBuildConsumeQueueConcurrently) { - this.enableBuildConsumeQueueConcurrently = enableBuildConsumeQueueConcurrently; - } } diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index 727bc5c00e..ba4b530641 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -234,7 +234,7 @@ public class DefaultMessageStore implements MessageStore { } } - if (!brokerConfig.isEnableBuildConsumeQueueConcurrently()) { + if (!messageStoreConfig.isEnableBuildConsumeQueueConcurrently()) { this.reputMessageService = new ReputMessageService(); } else { this.reputMessageService = new ConcurrentReputMessageService(); @@ -689,7 +689,7 @@ public class DefaultMessageStore implements MessageStore { this.recoverTopicQueueTable(); - if (!brokerConfig.isEnableBuildConsumeQueueConcurrently()) { + if (!messageStoreConfig.isEnableBuildConsumeQueueConcurrently()) { this.reputMessageService = new ReputMessageService(); } else { this.reputMessageService = new ConcurrentReputMessageService();