From ffb4a3664e9d3d8e3ed5795fe0f84a2664a6fdb0 Mon Sep 17 00:00:00 2001 From: Aaron Ai Date: Thu, 24 Nov 2022 19:32:07 +0800 Subject: [PATCH] Apply MDC in broker container mode (#5587) --- .../rocketmq/broker/BrokerController.java | 6 +- .../broker/BrokerPreOnlineService.java | 2 +- .../DefaultConsumerIdsChangeListener.java | 2 +- .../broker/filtersrv/FilterServerManager.java | 2 +- .../broker/latency/BrokerFastFailure.java | 2 +- .../LmqPullRequestHoldService.java | 2 +- .../longpolling/PullRequestHoldService.java | 2 +- .../rocketmq/broker/out/BrokerOuterAPI.java | 10 +- .../processor/NotificationProcessor.java | 2 +- .../processor/PopBufferMergeService.java | 2 +- .../broker/processor/PopMessageProcessor.java | 4 +- .../broker/processor/PopReviveService.java | 2 +- .../topic/TopicQueueMappingCleanService.java | 2 +- .../TransactionalMessageCheckService.java | 2 +- .../src/main/resources/rmq.broker.logback.xml | 669 +++++++++++------- .../common/AbstractBrokerRunnable.java | 18 +- .../rocketmq/common/BrokerIdentity.java | 8 +- .../rocketmq/common/ThreadFactoryImpl.java | 2 +- .../rocketmq/container/BrokerContainer.java | 6 +- .../container/BrokerContainerStartup.java | 2 +- .../container/InnerBrokerController.java | 4 +- docs/cn/BrokerContainer.md | 31 +- .../store/AllocateMappedFileService.java | 2 +- .../org/apache/rocketmq/store/CommitLog.java | 8 +- .../rocketmq/store/DefaultMessageStore.java | 18 +- .../rocketmq/store/StoreStatsService.java | 2 +- .../rocketmq/store/ha/DefaultHAClient.java | 2 +- .../store/ha/DefaultHAConnection.java | 4 +- .../rocketmq/store/ha/DefaultHAService.java | 2 +- .../store/ha/GroupTransferService.java | 2 +- .../HAConnectionStateNotificationService.java | 2 +- .../ha/autoswitch/AutoSwitchHAClient.java | 2 +- .../ha/autoswitch/AutoSwitchHAConnection.java | 4 +- .../ha/autoswitch/AutoSwitchHAService.java | 2 +- .../rocketmq/store/index/IndexService.java | 2 +- .../rocketmq/store/kv/CompactionService.java | 2 +- .../store/timer/TimerMessageStore.java | 14 +- .../rocketmq/store/ha/HAServerTest.java | 2 +- 38 files changed, 478 insertions(+), 374 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index 2cb1c3bb52..5697afce3f 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -1562,7 +1562,7 @@ public class BrokerController { scheduledFutures.add(this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { try { if (System.currentTimeMillis() < shouldStartTime) { BrokerController.LOG.info("Register to namesrv after {}", shouldStartTime); @@ -1584,7 +1584,7 @@ public class BrokerController { scheduledFutures.add(this.syncBrokerMemberGroupExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { try { BrokerController.this.syncBrokerMemberGroup(); } catch (Throwable e) { @@ -1606,7 +1606,7 @@ public class BrokerController { protected void scheduleSendHeartbeat() { scheduledFutures.add(this.brokerHeartbeatExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { if (isIsolated) { return; } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java index ab5d39e8fd..de2ccb2939 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java @@ -53,7 +53,7 @@ public class BrokerPreOnlineService extends ServiceThread { @Override public String getServiceName() { if (this.brokerController != null && this.brokerController.getBrokerConfig().isInBrokerContainer()) { - return brokerController.getBrokerIdentity().getLoggerIdentifier() + BrokerPreOnlineService.class.getSimpleName(); + return brokerController.getBrokerIdentity().getIdentifier() + BrokerPreOnlineService.class.getSimpleName(); } return BrokerPreOnlineService.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java b/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java index 2dc6ebe917..2ce036a0ff 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java @@ -47,7 +47,7 @@ public class DefaultConsumerIdsChangeListener implements ConsumerIdsChangeListen scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(brokerController.getBrokerConfig()) { @Override - public void run2() { + public void run0() { try { notifyConsumerChange(); } catch (Exception e) { diff --git a/broker/src/main/java/org/apache/rocketmq/broker/filtersrv/FilterServerManager.java b/broker/src/main/java/org/apache/rocketmq/broker/filtersrv/FilterServerManager.java index 96e69ca08e..f6628a158f 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/filtersrv/FilterServerManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/filtersrv/FilterServerManager.java @@ -56,7 +56,7 @@ public class FilterServerManager { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(brokerController.getBrokerConfig()) { @Override - public void run2() { + public void run0() { try { FilterServerManager.this.createFilterServer(); } catch (Exception e) { diff --git a/broker/src/main/java/org/apache/rocketmq/broker/latency/BrokerFastFailure.java b/broker/src/main/java/org/apache/rocketmq/broker/latency/BrokerFastFailure.java index c3a160b2fd..d3d0bc8ba3 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/latency/BrokerFastFailure.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/latency/BrokerFastFailure.java @@ -64,7 +64,7 @@ public class BrokerFastFailure { public void start() { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.brokerController.getBrokerConfig()) { @Override - public void run2() { + public void run0() { if (brokerController.getBrokerConfig().isBrokerFastFailureEnable()) { cleanExpiredRequest(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/longpolling/LmqPullRequestHoldService.java b/broker/src/main/java/org/apache/rocketmq/broker/longpolling/LmqPullRequestHoldService.java index 43b273d01e..88e74fd6e5 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/longpolling/LmqPullRequestHoldService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/longpolling/LmqPullRequestHoldService.java @@ -33,7 +33,7 @@ public class LmqPullRequestHoldService extends PullRequestHoldService { @Override public String getServiceName() { if (brokerController != null && brokerController.getBrokerConfig().isInBrokerContainer()) { - return this.brokerController.getBrokerIdentity().getLoggerIdentifier() + LmqPullRequestHoldService.class.getSimpleName(); + return this.brokerController.getBrokerIdentity().getIdentifier() + LmqPullRequestHoldService.class.getSimpleName(); } return LmqPullRequestHoldService.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/longpolling/PullRequestHoldService.java b/broker/src/main/java/org/apache/rocketmq/broker/longpolling/PullRequestHoldService.java index 88ccbbcefd..e8da9d0c47 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/longpolling/PullRequestHoldService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/longpolling/PullRequestHoldService.java @@ -92,7 +92,7 @@ public class PullRequestHoldService extends ServiceThread { @Override public String getServiceName() { if (brokerController != null && brokerController.getBrokerConfig().isInBrokerContainer()) { - return this.brokerController.getBrokerIdentity().getLoggerIdentifier() + PullRequestHoldService.class.getSimpleName(); + return this.brokerController.getBrokerIdentity().getIdentifier() + PullRequestHoldService.class.getSimpleName(); } return PullRequestHoldService.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java index a6853350e5..8d16908743 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java @@ -314,7 +314,7 @@ public class BrokerOuterAPI { brokerOuterExecutor.execute(new AbstractBrokerRunnable(new BrokerIdentity(clusterName, brokerName, brokerId, isInBrokerContainer)) { @Override - public void run2() { + public void run0() { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.QUERY_DATA_VERSION, requestHeader); request.setBody(dataVersion.encode()); @@ -346,7 +346,7 @@ public class BrokerOuterAPI { for (final String namesrvAddr : nameServerAddressList) { brokerOuterExecutor.execute(new AbstractBrokerRunnable(new BrokerIdentity(clusterName, brokerName, brokerId, isInBrokerContainer)) { @Override - public void run2() { + public void run0() { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.BROKER_HEARTBEAT, requestHeader); try { @@ -489,7 +489,7 @@ public class BrokerOuterAPI { for (final String namesrvAddr : nameServerAddressList) { brokerOuterExecutor.execute(new AbstractBrokerRunnable(brokerIdentity) { @Override - public void run2() { + public void run0() { try { RegisterBrokerResult result = registerBroker(namesrvAddr, oneway, timeoutMills, requestHeader, body); if (result != null) { @@ -622,7 +622,7 @@ public class BrokerOuterAPI { for (final String namesrvAddr : nameServerAddressList) { brokerOuterExecutor.execute(new AbstractBrokerRunnable(new BrokerIdentity(clusterName, brokerName, brokerId, isInBrokerContainer)) { @Override - public void run2() { + public void run0() { try { QueryDataVersionRequestHeader requestHeader = new QueryDataVersionRequestHeader(); requestHeader.setBrokerAddr(brokerAddr); @@ -1230,7 +1230,7 @@ public class BrokerOuterAPI { requestHeader.setConfirmOffset(confirmOffset); brokerOuterExecutor.execute(new AbstractBrokerRunnable(new BrokerIdentity(clusterName, brokerName, brokerId, isInBrokerContainer)) { @Override - public void run2() { + public void run0() { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.BROKER_HEARTBEAT, requestHeader); try { diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/NotificationProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/NotificationProcessor.java index eab7906e67..0b580df0fa 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/NotificationProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/NotificationProcessor.java @@ -57,7 +57,7 @@ public class NotificationProcessor implements NettyRequestProcessor { this.brokerController = brokerController; this.checkNotificationPollingThread = new Thread(new AbstractBrokerRunnable(brokerController.getBrokerConfig()) { @Override - public void run2() { + public void run0() { while (true) { if (Thread.currentThread().isInterrupted()) { break; diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java index da820a9d1c..2ccf4b8b30 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java @@ -77,7 +77,7 @@ public class PopBufferMergeService extends ServiceThread { @Override public String getServiceName() { if (this.brokerController != null && this.brokerController.getBrokerConfig().isInBrokerContainer()) { - return brokerController.getBrokerIdentity().getLoggerIdentifier() + PopBufferMergeService.class.getSimpleName(); + return brokerController.getBrokerIdentity().getIdentifier() + PopBufferMergeService.class.getSimpleName(); } return PopBufferMergeService.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java index 25beddb6d3..dba56102a2 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java @@ -783,7 +783,7 @@ public class PopMessageProcessor implements NettyRequestProcessor { @Override public String getServiceName() { if (PopMessageProcessor.this.brokerController.getBrokerConfig().isInBrokerContainer()) { - return PopMessageProcessor.this.brokerController.getBrokerIdentity().getLoggerIdentifier() + PopLongPollingService.class.getSimpleName(); + return PopMessageProcessor.this.brokerController.getBrokerIdentity().getIdentifier() + PopLongPollingService.class.getSimpleName(); } return PopLongPollingService.class.getSimpleName(); } @@ -1026,7 +1026,7 @@ public class PopMessageProcessor implements NettyRequestProcessor { @Override public String getServiceName() { if (PopMessageProcessor.this.brokerController.getBrokerConfig().isInBrokerContainer()) { - return PopMessageProcessor.this.brokerController.getBrokerIdentity().getLoggerIdentifier() + QueueLockManager.class.getSimpleName(); + return PopMessageProcessor.this.brokerController.getBrokerIdentity().getIdentifier() + QueueLockManager.class.getSimpleName(); } return QueueLockManager.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java index 173a3a8f3a..49de7433b4 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java @@ -72,7 +72,7 @@ public class PopReviveService extends ServiceThread { @Override public String getServiceName() { if (brokerController != null && brokerController.getBrokerConfig().isInBrokerContainer()) { - return brokerController.getBrokerIdentity().getLoggerIdentifier() + "PopReviveService_" + this.queueId; + return brokerController.getBrokerIdentity().getIdentifier() + "PopReviveService_" + this.queueId; } return "PopReviveService_" + this.queueId; } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicQueueMappingCleanService.java b/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicQueueMappingCleanService.java index f0e555565d..7047ef8b47 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicQueueMappingCleanService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicQueueMappingCleanService.java @@ -70,7 +70,7 @@ public class TopicQueueMappingCleanService extends ServiceThread { @Override public String getServiceName() { if (this.brokerConfig.isInBrokerContainer()) { - return this.brokerController.getBrokerIdentity().getLoggerIdentifier() + TopicQueueMappingCleanService.class.getSimpleName(); + return this.brokerController.getBrokerIdentity().getIdentifier() + TopicQueueMappingCleanService.class.getSimpleName(); } return TopicQueueMappingCleanService.class.getSimpleName(); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java b/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java index 9f4fef20de..52209c3fbd 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java @@ -34,7 +34,7 @@ public class TransactionalMessageCheckService extends ServiceThread { @Override public String getServiceName() { if (brokerController != null && brokerController.getBrokerConfig().isInBrokerContainer()) { - return brokerController.getBrokerIdentity().getLoggerIdentifier() + TransactionalMessageCheckService.class.getSimpleName(); + return brokerController.getBrokerIdentity().getIdentifier() + TransactionalMessageCheckService.class.getSimpleName(); } return TransactionalMessageCheckService.class.getSimpleName(); } diff --git a/broker/src/main/resources/rmq.broker.logback.xml b/broker/src/main/resources/rmq.broker.logback.xml index 10e6d3e5d0..9ba7054a9f 100644 --- a/broker/src/main/resources/rmq.broker.logback.xml +++ b/broker/src/main/resources/rmq.broker.logback.xml @@ -17,289 +17,410 @@ --> - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/broker_default.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/broker_default.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - + + + + brokerContainerLogDir + ${file.separator} + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}broker_default.log + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}broker_default.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/broker.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/broker.%i.log.gz - 1 - 20 - - - 128MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}broker.log + true + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}broker.%i.log.gz + + 1 + 20 + + + 128MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/protection.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/protection.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}protection.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}protection.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/watermark.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/watermark.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}watermark.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}watermark.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/store.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/store.%i.log.gz - 1 - 10 - - - 128MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}store.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}store.%i.log.gz + + 1 + 10 + + + 128MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/broker_traffic.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/broker_traffic.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + ${user.home}/logs/rocketmqlogs/broker_traffic.log + true + + ${user.home}/logs/rocketmqlogs/otherdays/broker_traffic.%i.log.gz + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/remoting.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/remoting.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}remoting.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}remoting.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/storeerror.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/storeerror.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}storeerror.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}storeerror.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/transaction.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/transaction.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}transaction.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}transaction.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/lock.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/lock.%i.log.gz - 1 - 5 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}lock.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}lock.%i.log.gz + + 1 + 5 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/filter.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/filter.%i.log.gz - 1 - 10 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}filter.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}filter.%i.log.gz + + 1 + 10 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/stats.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/stats.%i.log.gz - 1 - 5 - - - 100MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p - %m%n - UTF-8 - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}stats.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}stats.%i.log.gz + + 1 + 5 + + + 100MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p - %m%n + UTF-8 + + + - - ${user.home}/logs/rocketmqlogs/${brokerLogDir}/commercial.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/${brokerLogDir}/commercial.%i.log.gz - 1 - 10 - - - 500MB - + + + brokerContainerLogDir + ${file.separator} + + + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}commercial.log + + true + + + ${user.home}${file.separator}logs${file.separator}rocketmqlogs${brokerLogDir:-${file.separator}}${brokerContainerLogDir}${file.separator}otherdays${file.separator}commercial.%i.log.gz + + 1 + 10 + + + 500MB + + + - - ${user.home}/logs/rocketmqlogs/pop.log - true - - ${user.home}/logs/rocketmqlogs/otherdays/pop.%i.log - - 1 - 20 - - - 128MB - - - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n - UTF-8 - - - - - + + + brokerContainerLogDir + ${file.separator} + + + + ${user.home}/logs/rocketmqlogs/pop.log + true + + ${user.home}/logs/rocketmqlogs/otherdays/pop.%i.log + + 1 + 20 + + + 128MB + + + %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + UTF-8 + + + @@ -312,62 +433,62 @@ - + - + - + - + - + - + - + - + - + - + - + - + @@ -376,17 +497,17 @@ - - + + - - + + - + diff --git a/common/src/main/java/org/apache/rocketmq/common/AbstractBrokerRunnable.java b/common/src/main/java/org/apache/rocketmq/common/AbstractBrokerRunnable.java index 8bce59d762..34aabc5772 100644 --- a/common/src/main/java/org/apache/rocketmq/common/AbstractBrokerRunnable.java +++ b/common/src/main/java/org/apache/rocketmq/common/AbstractBrokerRunnable.java @@ -17,6 +17,9 @@ package org.apache.rocketmq.common; +import java.io.File; +import org.apache.rocketmq.logging.org.slf4j.MDC; + public abstract class AbstractBrokerRunnable implements Runnable { protected final BrokerIdentity brokerIdentity; @@ -24,17 +27,22 @@ public abstract class AbstractBrokerRunnable implements Runnable { this.brokerIdentity = brokerIdentity; } + private static final String MDC_BROKER_CONTAINER_LOG_DIR = "brokerContainerLogDir"; + /** * real logic for running */ - public abstract void run2(); + public abstract void run0(); @Override public void run() { - if (brokerIdentity.isInBrokerContainer()) { - // set threadlocal broker identity to forward logging to corresponding broker -// InnerLoggerFactory.BROKER_IDENTITY.set(brokerIdentity.getCanonicalName()); + try { + if (brokerIdentity.isInBrokerContainer()) { + MDC.put(MDC_BROKER_CONTAINER_LOG_DIR, File.separator + brokerIdentity.getCanonicalName()); + } + run0(); + } finally { + MDC.clear(); } - run2(); } } diff --git a/common/src/main/java/org/apache/rocketmq/common/BrokerIdentity.java b/common/src/main/java/org/apache/rocketmq/common/BrokerIdentity.java index 28e867ad54..4115744a42 100644 --- a/common/src/main/java/org/apache/rocketmq/common/BrokerIdentity.java +++ b/common/src/main/java/org/apache/rocketmq/common/BrokerIdentity.java @@ -115,13 +115,11 @@ public class BrokerIdentity { } public String getCanonicalName() { - if (isBrokerContainer) { -// return InnerLoggerFactory.BROKER_CONTAINER_NAME; - } - return this.getBrokerClusterName() + "_" + this.getBrokerName() + "_" + this.getBrokerId(); + return isBrokerContainer ? "BrokerContainer" : String.format("%s_%s_%d", brokerClusterName, brokerName, + brokerId); } - public String getLoggerIdentifier() { + public String getIdentifier() { return "#" + getCanonicalName() + "#"; } diff --git a/common/src/main/java/org/apache/rocketmq/common/ThreadFactoryImpl.java b/common/src/main/java/org/apache/rocketmq/common/ThreadFactoryImpl.java index cb6d0d71ae..bb0d141da0 100644 --- a/common/src/main/java/org/apache/rocketmq/common/ThreadFactoryImpl.java +++ b/common/src/main/java/org/apache/rocketmq/common/ThreadFactoryImpl.java @@ -47,7 +47,7 @@ public class ThreadFactoryImpl implements ThreadFactory { public ThreadFactoryImpl(final String threadNamePrefix, boolean daemon, BrokerIdentity brokerIdentity) { this.daemon = daemon; if (brokerIdentity != null && brokerIdentity.isInBrokerContainer()) { - this.threadNamePrefix = brokerIdentity.getLoggerIdentifier() + threadNamePrefix; + this.threadNamePrefix = brokerIdentity.getIdentifier() + threadNamePrefix; } else { this.threadNamePrefix = threadNamePrefix; } diff --git a/container/src/main/java/org/apache/rocketmq/container/BrokerContainer.java b/container/src/main/java/org/apache/rocketmq/container/BrokerContainer.java index 680f384439..47ca8e7041 100644 --- a/container/src/main/java/org/apache/rocketmq/container/BrokerContainer.java +++ b/container/src/main/java/org/apache/rocketmq/container/BrokerContainer.java @@ -158,7 +158,7 @@ public class BrokerContainer implements IBrokerContainer { // also auto update namesrv if specify this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(BrokerIdentity.BROKER_CONTAINER_IDENTITY) { @Override - public void run2() { + public void run0() { try { BrokerContainer.this.updateNamesrvAddr(); } catch (Throwable e) { @@ -170,7 +170,7 @@ public class BrokerContainer implements IBrokerContainer { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(BrokerIdentity.BROKER_CONTAINER_IDENTITY) { @Override - public void run2() { + public void run0() { try { BrokerContainer.this.brokerOuterAPI.fetchNameServerAddr(); } catch (Throwable e) { @@ -182,7 +182,7 @@ public class BrokerContainer implements IBrokerContainer { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(BrokerIdentity.BROKER_CONTAINER_IDENTITY) { @Override - public void run2() { + public void run0() { try { BrokerContainer.this.brokerOuterAPI.refreshMetadata(); } catch (Exception e) { diff --git a/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java b/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java index 02db476a05..f909e623b2 100644 --- a/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java +++ b/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java @@ -168,7 +168,7 @@ public class BrokerContainerStartup { brokerConfig.setBrokerConfigPath(filePath); - log = LoggerFactory.getLogger(brokerConfig.getLoggerIdentifier() + LoggerName.BROKER_LOGGER_NAME); + log = LoggerFactory.getLogger(brokerConfig.getIdentifier() + LoggerName.BROKER_LOGGER_NAME); MixAll.printObjectProperties(log, brokerConfig); MixAll.printObjectProperties(log, messageStoreConfig); diff --git a/container/src/main/java/org/apache/rocketmq/container/InnerBrokerController.java b/container/src/main/java/org/apache/rocketmq/container/InnerBrokerController.java index 47edd56e59..a1c1eecf59 100644 --- a/container/src/main/java/org/apache/rocketmq/container/InnerBrokerController.java +++ b/container/src/main/java/org/apache/rocketmq/container/InnerBrokerController.java @@ -69,7 +69,7 @@ public class InnerBrokerController extends BrokerController { scheduledFutures.add(this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { try { if (System.currentTimeMillis() < shouldStartTime) { BrokerController.LOG.info("Register to namesrv after {}", shouldStartTime); @@ -91,7 +91,7 @@ public class InnerBrokerController extends BrokerController { scheduledFutures.add(this.syncBrokerMemberGroupExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { try { InnerBrokerController.this.syncBrokerMemberGroup(); } catch (Throwable e) { diff --git a/docs/cn/BrokerContainer.md b/docs/cn/BrokerContainer.md index 8c52a677b6..236439284b 100644 --- a/docs/cn/BrokerContainer.md +++ b/docs/cn/BrokerContainer.md @@ -119,33 +119,10 @@ eg.假设broker获取到的配额是500g(根据replicasPerDiskPartition计算 ## 日志变化 -在BrokerContainer模式下并开启日志分离后,日志的默认输出路径将发生变化,每个broker日志的具体路径变化为 +在BrokerContainer模式下日志的默认输出路径将发生变化,具体为: + ``` -{user.home}/logs/{$brokerCanonicalName}_rocketmqlogs/ +{user.home}/logs/rocketmqlogs/${brokerCanonicalName}/ ``` -其中brokerCanonicalName为{BrokerClusterName_BrokerName_BrokerId},{BrokerClusterName_BrokerName_BrokerId}。 - -**开发者需要注意!** - -在BrokerContainer模式下,多个broker会在同一个BrokerContainer进程中,BrokerContainer模式下将提供broker日志分离功能,不同broker的日志将会输出到不同文件中。 - -主要通过线程名(ThreadName)或者通过设置线程本地变量(ThreadLocal)来区分不同broker线程,并且hack logback的logAppender将日志重定向到不同的文件中。 - -通过设置线程名来区分不同broker线程,线程名前缀必须是#BrokerClusterName_BrokerName_BrokerId# - -通过设置线程本地变量区分不同broker线程,设置的变量为BrokerClusterName_BrokerName_BrokerId -```java -// set threadlocal broker identity to forward logging to corresponding broker -InnerLoggerFactory.brokerIdentity.set(brokerIdentity.getCanonicalName()) -``` - -如果线程没有上述区分,日志将仍然会输出在原来的目录下。 - -以普通方式启动Broker(非BrokerContainer模式)时,日志将仍然会输出在原来的目录下。 - -具体实现方式可以参考Slf4jLoggerFactory和BrokerLogbackConfigurator两个类。 - -通过线程名和线程本地变量区分可以参考org.apache.rocketmq.common.AbstractBrokerRunnable、org.apache.rocketmq.common.ThreadFactoryImpl以及各个ServiceThread中getServiceName的实现。 - -参考文档:[原RIP](https://github.com/apache/rocketmq/wiki/RIP-31-Support-RocketMQ-BrokerContainer) \ No newline at end of file +其中 `brokerCanonicalName` 为 `{BrokerClusterName_BrokerName_BrokerId}`。 \ No newline at end of file diff --git a/store/src/main/java/org/apache/rocketmq/store/AllocateMappedFileService.java b/store/src/main/java/org/apache/rocketmq/store/AllocateMappedFileService.java index b2fcc2ceb5..4d2fc51683 100644 --- a/store/src/main/java/org/apache/rocketmq/store/AllocateMappedFileService.java +++ b/store/src/main/java/org/apache/rocketmq/store/AllocateMappedFileService.java @@ -122,7 +122,7 @@ public class AllocateMappedFileService extends ServiceThread { @Override public String getServiceName() { if (messageStore != null && messageStore.getBrokerConfig().isInBrokerContainer()) { - return messageStore.getBrokerIdentity().getLoggerIdentifier() + AllocateMappedFileService.class.getSimpleName(); + return messageStore.getBrokerIdentity().getIdentifier() + AllocateMappedFileService.class.getSimpleName(); } return AllocateMappedFileService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java index 2ea07695c9..d876f73b0b 100644 --- a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java +++ b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java @@ -1240,7 +1240,7 @@ public class CommitLog implements Swappable { @Override public String getServiceName() { if (CommitLog.this.defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return CommitLog.this.defaultMessageStore.getBrokerIdentity().getLoggerIdentifier() + CommitRealTimeService.class.getSimpleName(); + return CommitLog.this.defaultMessageStore.getBrokerIdentity().getIdentifier() + CommitRealTimeService.class.getSimpleName(); } return CommitRealTimeService.class.getSimpleName(); } @@ -1358,7 +1358,7 @@ public class CommitLog implements Swappable { @Override public String getServiceName() { if (CommitLog.this.defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return CommitLog.this.defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + FlushRealTimeService.class.getSimpleName(); + return CommitLog.this.defaultMessageStore.getBrokerConfig().getIdentifier() + FlushRealTimeService.class.getSimpleName(); } return FlushRealTimeService.class.getSimpleName(); } @@ -1503,7 +1503,7 @@ public class CommitLog implements Swappable { @Override public String getServiceName() { if (CommitLog.this.defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return CommitLog.this.defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + GroupCommitService.class.getSimpleName(); + return CommitLog.this.defaultMessageStore.getBrokerConfig().getIdentifier() + GroupCommitService.class.getSimpleName(); } return GroupCommitService.class.getSimpleName(); } @@ -1614,7 +1614,7 @@ public class CommitLog implements Swappable { @Override public String getServiceName() { if (CommitLog.this.defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return CommitLog.this.defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + GroupCheckService.class.getSimpleName(); + return CommitLog.this.defaultMessageStore.getBrokerConfig().getIdentifier() + GroupCheckService.class.getSimpleName(); } return GroupCheckService.class.getSimpleName(); } 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 60c9331d46..31c1a2cb41 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -1593,21 +1593,21 @@ public class DefaultMessageStore implements MessageStore { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { DefaultMessageStore.this.cleanFilesPeriodically(); } }, 1000 * 60, this.messageStoreConfig.getCleanResourceInterval(), TimeUnit.MILLISECONDS); this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { DefaultMessageStore.this.checkSelf(); } }, 1, 10, TimeUnit.MINUTES); this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { if (DefaultMessageStore.this.getMessageStoreConfig().isDebugLockEnable()) { try { if (DefaultMessageStore.this.commitLog.getBeginTimeInLock() != 0) { @@ -1628,7 +1628,7 @@ public class DefaultMessageStore implements MessageStore { this.scheduledExecutorService.scheduleAtFixedRate(new AbstractBrokerRunnable(this.getBrokerIdentity()) { @Override - public void run2() { + public void run0() { DefaultMessageStore.this.storeCheckpoint.flush(); } }, 1, 1, TimeUnit.SECONDS); @@ -2056,7 +2056,7 @@ public class DefaultMessageStore implements MessageStore { } public String getServiceName() { - return DefaultMessageStore.this.brokerConfig.getLoggerIdentifier() + CleanCommitLogService.class.getSimpleName(); + return DefaultMessageStore.this.brokerConfig.getIdentifier() + CleanCommitLogService.class.getSimpleName(); } private boolean isTimeToDelete() { @@ -2253,7 +2253,7 @@ public class DefaultMessageStore implements MessageStore { } public String getServiceName() { - return DefaultMessageStore.this.brokerConfig.getLoggerIdentifier() + CleanConsumeQueueService.class.getSimpleName(); + return DefaultMessageStore.this.brokerConfig.getIdentifier() + CleanConsumeQueueService.class.getSimpleName(); } } @@ -2371,7 +2371,7 @@ public class DefaultMessageStore implements MessageStore { public String getServiceName() { if (brokerConfig.isInBrokerContainer()) { - return brokerConfig.getLoggerIdentifier() + CorrectLogicOffsetService.class.getSimpleName(); + return brokerConfig.getIdentifier() + CorrectLogicOffsetService.class.getSimpleName(); } return CorrectLogicOffsetService.class.getSimpleName(); } @@ -2443,7 +2443,7 @@ public class DefaultMessageStore implements MessageStore { @Override public String getServiceName() { if (DefaultMessageStore.this.brokerConfig.isInBrokerContainer()) { - return DefaultMessageStore.this.getBrokerIdentity().getLoggerIdentifier() + FlushConsumeQueueService.class.getSimpleName(); + return DefaultMessageStore.this.getBrokerIdentity().getIdentifier() + FlushConsumeQueueService.class.getSimpleName(); } return FlushConsumeQueueService.class.getSimpleName(); } @@ -2619,7 +2619,7 @@ public class DefaultMessageStore implements MessageStore { @Override public String getServiceName() { if (DefaultMessageStore.this.getBrokerConfig().isInBrokerContainer()) { - return DefaultMessageStore.this.getBrokerIdentity().getLoggerIdentifier() + ReputMessageService.class.getSimpleName(); + return DefaultMessageStore.this.getBrokerIdentity().getIdentifier() + ReputMessageService.class.getSimpleName(); } return ReputMessageService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java b/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java index bb6dc9774f..1969b146aa 100644 --- a/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java +++ b/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java @@ -546,7 +546,7 @@ public class StoreStatsService extends ServiceThread { @Override public String getServiceName() { if (this.brokerIdentity != null && this.brokerIdentity.isInBrokerContainer()) { - return brokerIdentity.getLoggerIdentifier() + StoreStatsService.class.getSimpleName(); + return brokerIdentity.getIdentifier() + StoreStatsService.class.getSimpleName(); } return StoreStatsService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAClient.java b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAClient.java index fce9f08683..02668558a2 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAClient.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAClient.java @@ -388,7 +388,7 @@ public class DefaultHAClient extends ServiceThread implements HAClient { @Override public String getServiceName() { if (this.defaultMessageStore != null && this.defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return this.defaultMessageStore.getBrokerIdentity().getLoggerIdentifier() + DefaultHAClient.class.getSimpleName(); + return this.defaultMessageStore.getBrokerIdentity().getIdentifier() + DefaultHAClient.class.getSimpleName(); } return DefaultHAClient.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAConnection.java b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAConnection.java index de7bfe3d77..8b3598666f 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAConnection.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAConnection.java @@ -185,7 +185,7 @@ public class DefaultHAConnection implements HAConnection { @Override public String getServiceName() { if (haService.getDefaultMessageStore().getBrokerConfig().isInBrokerContainer()) { - return haService.getDefaultMessageStore().getBrokerIdentity().getLoggerIdentifier() + ReadSocketService.class.getSimpleName(); + return haService.getDefaultMessageStore().getBrokerIdentity().getIdentifier() + ReadSocketService.class.getSimpleName(); } return ReadSocketService.class.getSimpleName(); } @@ -440,7 +440,7 @@ public class DefaultHAConnection implements HAConnection { @Override public String getServiceName() { if (haService.getDefaultMessageStore().getBrokerConfig().isInBrokerContainer()) { - return haService.getDefaultMessageStore().getBrokerIdentity().getLoggerIdentifier() + WriteSocketService.class.getSimpleName(); + return haService.getDefaultMessageStore().getBrokerIdentity().getIdentifier() + WriteSocketService.class.getSimpleName(); } return WriteSocketService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java index a9f8a38358..98deb233f3 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java @@ -270,7 +270,7 @@ public class DefaultHAService implements HAService { @Override public String getServiceName() { if (defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + AcceptSocketService.class.getSimpleName(); + return defaultMessageStore.getBrokerConfig().getIdentifier() + AcceptSocketService.class.getSimpleName(); } return DefaultAcceptSocketService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java b/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java index 8896e74986..5318dee8f5 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java @@ -169,7 +169,7 @@ public class GroupTransferService extends ServiceThread { @Override public String getServiceName() { if (defaultMessageStore != null && defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return defaultMessageStore.getBrokerIdentity().getLoggerIdentifier() + GroupTransferService.class.getSimpleName(); + return defaultMessageStore.getBrokerIdentity().getIdentifier() + GroupTransferService.class.getSimpleName(); } return GroupTransferService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/HAConnectionStateNotificationService.java b/store/src/main/java/org/apache/rocketmq/store/ha/HAConnectionStateNotificationService.java index 750f1ca4de..197d9f6ba4 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/HAConnectionStateNotificationService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/HAConnectionStateNotificationService.java @@ -47,7 +47,7 @@ public class HAConnectionStateNotificationService extends ServiceThread { @Override public String getServiceName() { if (defaultMessageStore != null && defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return defaultMessageStore.getBrokerIdentity().getLoggerIdentifier() + HAConnectionStateNotificationService.class.getSimpleName(); + return defaultMessageStore.getBrokerIdentity().getIdentifier() + HAConnectionStateNotificationService.class.getSimpleName(); } return HAConnectionStateNotificationService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java index 55ca1cc170..4e0e37aed4 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java @@ -136,7 +136,7 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient { @Override public String getServiceName() { if (haService.getDefaultMessageStore().getBrokerConfig().isInBrokerContainer()) { - return haService.getDefaultMessageStore().getBrokerIdentity().getLoggerIdentifier() + AutoSwitchHAClient.class.getSimpleName(); + return haService.getDefaultMessageStore().getBrokerIdentity().getIdentifier() + AutoSwitchHAClient.class.getSimpleName(); } return AutoSwitchHAClient.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java index 3e7d0cb815..1afb9f6dec 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java @@ -300,7 +300,7 @@ public class AutoSwitchHAConnection implements HAConnection { @Override public String getServiceName() { if (haService.getDefaultMessageStore().getBrokerConfig().isInBrokerContainer()) { - return haService.getDefaultMessageStore().getBrokerIdentity().getLoggerIdentifier() + ReadSocketService.class.getSimpleName(); + return haService.getDefaultMessageStore().getBrokerIdentity().getIdentifier() + ReadSocketService.class.getSimpleName(); } return ReadSocketService.class.getSimpleName(); } @@ -446,7 +446,7 @@ public class AutoSwitchHAConnection implements HAConnection { @Override public String getServiceName() { if (haService.getDefaultMessageStore().getBrokerConfig().isInBrokerContainer()) { - return haService.getDefaultMessageStore().getBrokerIdentity().getLoggerIdentifier() + WriteSocketService.class.getSimpleName(); + return haService.getDefaultMessageStore().getBrokerIdentity().getIdentifier() + WriteSocketService.class.getSimpleName(); } return WriteSocketService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java index 72612d3080..59eb140333 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java @@ -409,7 +409,7 @@ public class AutoSwitchHAService extends DefaultHAService { @Override public String getServiceName() { if (defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + AcceptSocketService.class.getSimpleName(); + return defaultMessageStore.getBrokerConfig().getIdentifier() + AcceptSocketService.class.getSimpleName(); } return AutoSwitchAcceptSocketService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/index/IndexService.java b/store/src/main/java/org/apache/rocketmq/store/index/IndexService.java index 89a69e375b..ef5d21ac7c 100644 --- a/store/src/main/java/org/apache/rocketmq/store/index/IndexService.java +++ b/store/src/main/java/org/apache/rocketmq/store/index/IndexService.java @@ -342,7 +342,7 @@ public class IndexService { Thread flushThread = new Thread(new AbstractBrokerRunnable(defaultMessageStore.getBrokerConfig()) { @Override - public void run2() { + public void run0() { IndexService.this.flush(flushThisFile); } }, "FlushIndexFileThread"); diff --git a/store/src/main/java/org/apache/rocketmq/store/kv/CompactionService.java b/store/src/main/java/org/apache/rocketmq/store/kv/CompactionService.java index 3688ab6416..1b5d389135 100644 --- a/store/src/main/java/org/apache/rocketmq/store/kv/CompactionService.java +++ b/store/src/main/java/org/apache/rocketmq/store/kv/CompactionService.java @@ -74,7 +74,7 @@ public class CompactionService extends ServiceThread { @Override public String getServiceName() { if (defaultMessageStore != null && defaultMessageStore.getBrokerConfig().isInBrokerContainer()) { - return defaultMessageStore.getBrokerConfig().getLoggerIdentifier() + CompactionService.class.getSimpleName(); + return defaultMessageStore.getBrokerConfig().getIdentifier() + CompactionService.class.getSimpleName(); } return CompactionService.class.getSimpleName(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java index 1c7253c823..89b93abd0a 100644 --- a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMessageStore.java @@ -1235,7 +1235,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1261,7 +1261,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1342,7 +1342,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1386,7 +1386,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1454,7 +1454,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1537,7 +1537,7 @@ public class TimerMessageStore { public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } @@ -1572,7 +1572,7 @@ public class TimerMessageStore { @Override public String getServiceName() { String brokerIdentifier = ""; if (TimerMessageStore.this.messageStore instanceof DefaultMessageStore && ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().isInBrokerContainer()) { - brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getLoggerIdentifier(); + brokerIdentifier = ((DefaultMessageStore) TimerMessageStore.this.messageStore).getBrokerConfig().getIdentifier(); } return brokerIdentifier + this.getClass().getSimpleName(); } diff --git a/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java b/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java index a8ce5179dc..54174ac166 100644 --- a/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java @@ -261,7 +261,7 @@ public class HAServerTest { BrokerConfig brokerConfig = mock(BrokerConfig.class); doReturn(true).when(brokerConfig).isInBrokerContainer(); - doReturn("mock").when(brokerConfig).getLoggerIdentifier(); + doReturn("mock").when(brokerConfig).getIdentifier(); doReturn(brokerConfig).when(messageStore).getBrokerConfig(); doReturn(new SystemClock()).when(messageStore).getSystemClock(); doAnswer(invocation -> System.currentTimeMillis()).when(messageStore).now();