mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 05:44:03 +08:00
@@ -136,7 +136,7 @@ public class PlainPermissionManager {
|
||||
log.warn("No data in file {}", currentFile);
|
||||
continue;
|
||||
}
|
||||
log.info("Broker plain acl conf data is : ", plainAclConfData.toString());
|
||||
log.info("Broker plain acl conf data is : {}", plainAclConfData.toString());
|
||||
|
||||
List<RemoteAddressStrategy> globalWhiteRemoteAddressStrategyList = new ArrayList<>();
|
||||
JSONArray globalWhiteRemoteAddressesList = plainAclConfData.getJSONArray("globalWhiteRemoteAddresses");
|
||||
@@ -218,7 +218,7 @@ public class PlainPermissionManager {
|
||||
log.warn("No data in {}, skip it", aclFilePath);
|
||||
return;
|
||||
}
|
||||
log.info("Broker plain acl conf data is : ", plainAclConfData.toString());
|
||||
log.info("Broker plain acl conf data is : {}", plainAclConfData.toString());
|
||||
JSONArray globalWhiteRemoteAddressesList = plainAclConfData.getJSONArray("globalWhiteRemoteAddresses");
|
||||
if (globalWhiteRemoteAddressesList != null && !globalWhiteRemoteAddressesList.isEmpty()) {
|
||||
for (int i = 0; i < globalWhiteRemoteAddressesList.size(); i++) {
|
||||
|
||||
@@ -114,7 +114,6 @@ import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingTimeoutException;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingSerializable;
|
||||
@@ -142,7 +141,7 @@ import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import org.apache.rocketmq.store.config.BrokerRole;
|
||||
|
||||
public class AdminBrokerProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class AdminBrokerProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
private final BrokerController brokerController;
|
||||
|
||||
|
||||
+1
-2
@@ -39,11 +39,10 @@ import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class ClientManageProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class ClientManageProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
private final BrokerController brokerController;
|
||||
|
||||
|
||||
+1
-2
@@ -35,10 +35,9 @@ import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class ConsumerManageProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class ConsumerManageProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
|
||||
private final BrokerController brokerController;
|
||||
|
||||
+1
-2
@@ -33,7 +33,6 @@ import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.MessageExtBrokerInner;
|
||||
import org.apache.rocketmq.store.PutMessageResult;
|
||||
@@ -42,7 +41,7 @@ import org.apache.rocketmq.store.config.BrokerRole;
|
||||
/**
|
||||
* EndTransaction processor: process commit and rollback message
|
||||
*/
|
||||
public class EndTransactionProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class EndTransactionProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger LOGGER = InternalLoggerFactory.getLogger(LoggerName.TRANSACTION_LOGGER_NAME);
|
||||
private final BrokerController brokerController;
|
||||
|
||||
|
||||
+1
-2
@@ -22,10 +22,9 @@ import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class ForwardRequestProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class ForwardRequestProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
|
||||
private final BrokerController brokerController;
|
||||
|
||||
@@ -59,7 +59,6 @@ import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.common.RemotingUtil;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.RequestTask;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.GetMessageResult;
|
||||
@@ -69,7 +68,7 @@ import org.apache.rocketmq.store.PutMessageResult;
|
||||
import org.apache.rocketmq.store.config.BrokerRole;
|
||||
import org.apache.rocketmq.store.stats.BrokerStatsManager;
|
||||
|
||||
public class PullMessageProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class PullMessageProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
private final BrokerController brokerController;
|
||||
private List<ConsumeMessageHook> consumeMessageHookList;
|
||||
|
||||
+1
-2
@@ -34,12 +34,11 @@ import org.apache.rocketmq.common.protocol.header.QueryMessageResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ViewMessageRequestHeader;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.AsyncNettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.QueryMessageResult;
|
||||
import org.apache.rocketmq.store.SelectMappedBufferResult;
|
||||
|
||||
public class QueryMessageProcessor extends AsyncNettyRequestProcessor implements NettyRequestProcessor {
|
||||
public class QueryMessageProcessor extends AsyncNettyRequestProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
|
||||
private final BrokerController brokerController;
|
||||
|
||||
+1
-2
@@ -39,7 +39,6 @@ import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingException;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.MessageExtBrokerInner;
|
||||
import org.apache.rocketmq.store.PutMessageResult;
|
||||
@@ -48,7 +47,7 @@ import org.apache.rocketmq.store.stats.BrokerStatsManager;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.concurrent.ThreadLocalRandom;
|
||||
|
||||
public class ReplyMessageProcessor extends AbstractSendMessageProcessor implements NettyRequestProcessor {
|
||||
public class ReplyMessageProcessor extends AbstractSendMessageProcessor {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME);
|
||||
|
||||
public ReplyMessageProcessor(final BrokerController brokerController) {
|
||||
|
||||
@@ -51,7 +51,6 @@ import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
import org.apache.rocketmq.common.sysflag.TopicSysFlag;
|
||||
import org.apache.rocketmq.common.topic.TopicValidator;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.netty.RemotingResponseCallback;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.DefaultMessageStore;
|
||||
@@ -62,7 +61,7 @@ import org.apache.rocketmq.store.config.MessageStoreConfig;
|
||||
import org.apache.rocketmq.store.config.StorePathConfigHelper;
|
||||
import org.apache.rocketmq.store.stats.BrokerStatsManager;
|
||||
|
||||
public class SendMessageProcessor extends AbstractSendMessageProcessor implements NettyRequestProcessor {
|
||||
public class SendMessageProcessor extends AbstractSendMessageProcessor {
|
||||
|
||||
private List<ConsumeMessageHook> consumeMessageHookList;
|
||||
|
||||
|
||||
@@ -99,12 +99,12 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
this.traceTopicName = TopicValidator.RMQ_SYS_TRACE_TOPIC;
|
||||
}
|
||||
this.traceExecutor = new ThreadPoolExecutor(//
|
||||
10, //
|
||||
20, //
|
||||
1000 * 60, //
|
||||
TimeUnit.MILLISECONDS, //
|
||||
this.appenderQueue, //
|
||||
new ThreadFactoryImpl("MQTraceSendThread_"));
|
||||
10, //
|
||||
20, //
|
||||
1000 * 60, //
|
||||
TimeUnit.MILLISECONDS, //
|
||||
this.appenderQueue, //
|
||||
new ThreadFactoryImpl("MQTraceSendThread_"));
|
||||
traceProducer = getAndCreateTraceProducer(rpcHook);
|
||||
}
|
||||
|
||||
@@ -165,7 +165,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
traceProducerInstance.setSendMsgTimeout(5000);
|
||||
traceProducerInstance.setVipChannelEnabled(false);
|
||||
// The max size of message is 128K
|
||||
traceProducerInstance.setMaxMessageSize(maxMsgSize - 10 * 1000);
|
||||
traceProducerInstance.setMaxMessageSize(maxMsgSize);
|
||||
}
|
||||
return traceProducerInstance;
|
||||
}
|
||||
@@ -188,6 +188,11 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
// The maximum waiting time for refresh,avoid being written all the time, resulting in failure to return.
|
||||
long end = System.currentTimeMillis() + 500;
|
||||
while (System.currentTimeMillis() <= end) {
|
||||
synchronized (taskQueueByTopic) {
|
||||
for (TraceDataSegment taskInfo : taskQueueByTopic.values()) {
|
||||
taskInfo.sendAllData();
|
||||
}
|
||||
}
|
||||
synchronized (traceContextQueue) {
|
||||
if (traceContextQueue.size() == 0 && appenderQueue.size() == 0) {
|
||||
break;
|
||||
@@ -252,7 +257,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
while (System.currentTimeMillis() < endTime) {
|
||||
try {
|
||||
TraceContext traceContext = traceContextQueue.poll(
|
||||
endTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS
|
||||
endTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS
|
||||
);
|
||||
|
||||
if (traceContext != null && !traceContext.getTraceBeans().isEmpty()) {
|
||||
@@ -324,7 +329,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
initFirstBeanAddTime();
|
||||
this.traceTransferBeanList.add(traceTransferBean);
|
||||
this.currentMsgSize += traceTransferBean.getTransData().length();
|
||||
if (currentMsgSize >= traceProducer.getMaxMessageSize()) {
|
||||
if (currentMsgSize >= traceProducer.getMaxMessageSize() - 10 * 1000) {
|
||||
List<TraceTransferBean> dataToSend = new ArrayList(traceTransferBeanList);
|
||||
AsyncDataSendTask asyncDataSendTask = new AsyncDataSendTask(traceTopicName, regionId, dataToSend);
|
||||
traceExecutor.submit(asyncDataSendTask);
|
||||
|
||||
@@ -21,6 +21,8 @@ import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Random;
|
||||
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
|
||||
public class BrokerData implements Comparable<BrokerData> {
|
||||
@@ -96,12 +98,7 @@ public class BrokerData implements Comparable<BrokerData> {
|
||||
return false;
|
||||
} else if (!brokerAddrs.equals(other.brokerAddrs))
|
||||
return false;
|
||||
if (brokerName == null) {
|
||||
if (other.brokerName != null)
|
||||
return false;
|
||||
} else if (!brokerName.equals(other.brokerName))
|
||||
return false;
|
||||
return true;
|
||||
return StringUtils.equals(brokerName, other.brokerName);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
limitations under the License.
|
||||
-->
|
||||
|
||||
<configuration>
|
||||
<configuration scan="true" scanPeriod="30 seconds">
|
||||
<appender name="DefaultAppender"
|
||||
class="ch.qos.logback.core.rolling.RollingFileAppender">
|
||||
<file>${user.home}/logs/rocketmqlogs/${brokerLogDir}/broker_default.log</file>
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
limitations under the License.
|
||||
-->
|
||||
|
||||
<configuration>
|
||||
<configuration scan="true" scanPeriod="30 seconds">
|
||||
<appender name="DefaultAppender"
|
||||
class="ch.qos.logback.core.rolling.RollingFileAppender">
|
||||
<file>${user.home}/logs/rocketmqlogs/namesrv_default.log</file>
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
limitations under the License.
|
||||
-->
|
||||
|
||||
<configuration>
|
||||
<configuration scan="true" scanPeriod="30 seconds">
|
||||
<appender name="DefaultAppender"
|
||||
class="ch.qos.logback.core.rolling.RollingFileAppender">
|
||||
<file>${user.home}/logs/rocketmqlogs/tools_default.log</file>
|
||||
|
||||
+2
-2
@@ -107,7 +107,7 @@ RocketMQ分布式消息队列的消息过滤方式有别于其它MQ中间件,
|
||||
### 4 负载均衡
|
||||
RocketMQ中的负载均衡都在Client端完成,具体来说的话,主要可以分为Producer端发送消息时候的负载均衡和Consumer端订阅消息的负载均衡。
|
||||
#### 4.1 Producer的负载均衡
|
||||
Producer端在发送消息的时候,会先根据Topic找到指定的TopicPublishInfo,在获取了TopicPublishInfo路由信息后,RocketMQ的客户端在默认方式下selectOneMessageQueue()方法会从TopicPublishInfo中的messageQueueList中选择一个队列(MessageQueue)进行发送消息。具体的容错策略均在MQFaultStrategy这个类中定义。这里有一个sendLatencyFaultEnable开关变量,如果开启,在随机递增取模的基础上,再过滤掉not available的Broker代理。所谓的"latencyFaultTolerance",是指对之前失败的,按一定的时间做退避。例如,如果上次请求的latency超过550Lms,就退避3000Lms;超过1000L,就退避60000L;如果关闭,采用随机递增取模的方式选择一个队列(MessageQueue)来发送消息,latencyFaultTolerance机制是实现消息发送高可用的核心关键所在。
|
||||
Producer端在发送消息的时候,会先根据Topic找到指定的TopicPublishInfo,在获取了TopicPublishInfo路由信息后,RocketMQ的客户端在默认方式下selectOneMessageQueue()方法会从TopicPublishInfo中的messageQueueList中选择一个队列(MessageQueue)进行发送消息。具体的容错策略均在MQFaultStrategy这个类中定义。这里有一个sendLatencyFaultEnable开关变量,如果开启,在随机递增取模的基础上,再过滤掉not available的Broker代理。所谓的"latencyFaultTolerance",是指对之前失败的,按一定的时间做退避。例如,如果上次请求的latency超过550L ms,就退避30000L ms;超过1000L,就退避60000L;如果关闭,采用随机递增取模的方式选择一个队列(MessageQueue)来发送消息,latencyFaultTolerance机制是实现消息发送高可用的核心关键所在。
|
||||
#### 4.2 Consumer的负载均衡
|
||||
在RocketMQ中,Consumer端的两种消费模式(Push/Pull)都是基于拉模式来获取消息的,而在Push模式只是对pull模式的一种封装,其本质实现为消息拉取线程在从服务器拉取到一批消息后,然后提交到消息消费线程池后,又“马不停蹄”的继续向服务器再次尝试拉取消息。如果未拉取到消息,则延迟一下又继续拉取。在两种基于拉模式的消费方式(Push/Pull)中,均需要Consumer端知道从Broker端的哪一个消息队列中去获取消息。因此,有必要在Consumer端来做负载均衡,即Broker端中多个MessageQueue分配给同一个ConsumerGroup中的哪些Consumer消费。
|
||||
|
||||
@@ -137,7 +137,7 @@ Producer端在发送消息的时候,会先根据Topic找到指定的TopicPubli
|
||||
|
||||
- 上图中processQueueTable的绿色部分,表示与分配到的消息队列集合mqSet的交集。判断该ProcessQueue是否已经过期了,在Pull模式的不用管,如果是Push模式的,设置Dropped属性为true,并且调用removeUnnecessaryMessageQueue()方法,像上面一样尝试移除Entry;
|
||||
|
||||
最后,为过滤后的消息队列集合(mqSet)中的每个MessageQueue创建一个ProcessQueue对象并存入RebalanceImpl的processQueueTable队列中(其中调用RebalanceImpl实例的computePullFromWhere(MessageQueue mq)方法获取该MessageQueue对象的下一个进度消费值offset,随后填充至接下来要创建的pullRequest对象属性中),并创建拉取请求对象—pullRequest添加到拉取列表—pullRequestList中,最后执行dispatchPullRequest()方法,将Pull消息的请求对象PullRequest依次放入PullMessageService服务线程的阻塞队列pullRequestQueue中,待该服务线程取出后向Broker端发起Pull消息的请求。其中,可以重点对比下,RebalancePushImpl和RebalancePullImpl两个实现类的dispatchPullRequest()方法不同,RebalancePullImpl类里面的该方法为空,这样子也就回答了上一篇中最后的那道思考题了。
|
||||
最后,为过滤后的消息队列集合(mqSet)中的每个MessageQueue创建一个ProcessQueue对象并存入RebalanceImpl的processQueueTable队列中(其中调用RebalanceImpl实例的computePullFromWhere(MessageQueue mq)方法获取该MessageQueue对象的下一个进度消费值offset,随后填充至接下来要创建的pullRequest对象属性中),并创建拉取请求对象—pullRequest添加到拉取列表—pullRequestList中,最后执行dispatchPullRequest()方法,将Pull消息的请求对象PullRequest依次放入PullMessageService服务线程的阻塞队列pullRequestQueue中,待该服务线程取出后向Broker端发起Pull消息的请求。
|
||||
|
||||
消息消费队列在同一消费组不同消费者之间的负载均衡,其核心设计理念是在一个消息消费队列在同一时间只允许被同一消费组内的一个消费者消费,一个消息消费者能同时消费多个消息队列。
|
||||
|
||||
|
||||
@@ -3,7 +3,7 @@ Load balancing in RocketMQ is accomplished on Client side. Specifically, it can
|
||||
|
||||
### Producer Load Balancing
|
||||
When the Producer sends a message, it will first find the specified TopicPublishInfo according to Topic. After getting the routing information of TopicPublishInfo, the RocketMQ client will select a queue (MessageQueue) from the messageQueue List in TopicPublishInfo to send the message by default.Specific fault-tolerant strategies are defined in the MQFaultStrategy class.
|
||||
Here is a sendLatencyFaultEnable switch variable, which, if turned on, filters out the Broker agent of not available on the basis of randomly gradually increasing modular arithmetic selection. The so-called "latencyFault Tolerance" refers to a certain period of time to avoid previous failures. For example, if the latency of the last request exceeds 550 Lms, it will evade 3000 Lms; if it exceeds 1000L, it will evade 60000 L; if it is closed, it will choose a queue (MessageQueue) to send messages by randomly gradually increasing modular arithmetic, and the latencyFault Tolerance mechanism is the key to achieve high availability of message sending.
|
||||
Here is a sendLatencyFaultEnable switch variable, which, if turned on, filters out the Broker agent of not available on the basis of randomly gradually increasing modular arithmetic selection. The so-called "latencyFault Tolerance" refers to a certain period of time to avoid previous failures. For example, if the latency of the last request exceeds 550 Lms, it will evade 30000 Lms; if it exceeds 1000L, it will evade 60000L; if it is closed, it will choose a queue (MessageQueue) to send messages by randomly gradually increasing modular arithmetic, and the latencyFault Tolerance mechanism is the key to achieve high availability of message sending.
|
||||
|
||||
### Consumer Load Balancing
|
||||
In RocketMQ, the two consumption modes (Push/Pull) on the Consumer side are both based on the pull mode to get the message, while in the Push mode it is only a kind of encapsulation of the pull mode, which is essentially implemented as the message pulling thread after pulling a batch of messages from the server. After submitting to the message consuming thread pool, it continues to try again to pull the message to the server. If the message is not pulled, the pull is delayed and continues. In both pull mode based consumption patterns (Push/Pull), the Consumer needs to know which message queue - queue from the Broker side to get the message. Therefore, it is necessary to do load balancing on the Consumer side, that is, which Consumer consumption is allocated to the same ConsumerGroup by more than one MessageQueue on the Broker side.
|
||||
|
||||
@@ -604,6 +604,15 @@ public class CommitLog {
|
||||
return keyBuilder.toString();
|
||||
}
|
||||
|
||||
public void updateMaxMessageSize(PutMessageThreadLocal putMessageThreadLocal) {
|
||||
// dynamically adjust maxMessageSize, but not support DLedger mode temporarily
|
||||
int newMaxMessageSize = this.defaultMessageStore.getMessageStoreConfig().getMaxMessageSize();
|
||||
if (newMaxMessageSize >= 10 &&
|
||||
putMessageThreadLocal.getEncoder().getMaxMessageBodySize() != newMaxMessageSize) {
|
||||
putMessageThreadLocal.getEncoder().updateEncoderBufferCapacity(newMaxMessageSize);
|
||||
}
|
||||
}
|
||||
|
||||
public CompletableFuture<PutMessageResult> asyncPutMessage(final MessageExtBrokerInner msg) {
|
||||
// Set the storage time
|
||||
msg.setStoreTimestamp(System.currentTimeMillis());
|
||||
@@ -650,6 +659,7 @@ public class CommitLog {
|
||||
}
|
||||
|
||||
PutMessageThreadLocal putMessageThreadLocal = this.putMessageThreadLocal.get();
|
||||
updateMaxMessageSize(putMessageThreadLocal);
|
||||
if (!multiDispatch.isMultiDispatchMsg(msg)) {
|
||||
PutMessageResult encodeResult = putMessageThreadLocal.getEncoder().encode(msg);
|
||||
if (encodeResult != null) {
|
||||
@@ -768,6 +778,7 @@ public class CommitLog {
|
||||
|
||||
//fine-grained lock instead of the coarse-grained
|
||||
PutMessageThreadLocal pmThreadLocal = this.putMessageThreadLocal.get();
|
||||
updateMaxMessageSize(pmThreadLocal);
|
||||
MessageExtEncoder batchEncoder = pmThreadLocal.getEncoder();
|
||||
|
||||
PutMessageContext putMessageContext = new PutMessageContext(generateKey(pmThreadLocal.getKeyBuilder(), messageExtBatch));
|
||||
@@ -1479,15 +1490,16 @@ public class CommitLog {
|
||||
}
|
||||
|
||||
public static class MessageExtEncoder {
|
||||
private final ByteBuf byteBuf;
|
||||
private ByteBuf byteBuf;
|
||||
// The maximum length of the message body.
|
||||
private final int maxMessageBodySize;
|
||||
private int maxMessageBodySize;
|
||||
// The maximum length of the full message.
|
||||
private final int maxMessageSize;
|
||||
private int maxMessageSize;
|
||||
MessageExtEncoder(final int maxMessageBodySize) {
|
||||
ByteBufAllocator alloc = UnpooledByteBufAllocator.DEFAULT;
|
||||
//Reserve 64kb for encoding buffer outside body
|
||||
int maxMessageSize = maxMessageBodySize + 64 * 1024;
|
||||
int maxMessageSize = Integer.MAX_VALUE - maxMessageBodySize >= 64 * 1024 ?
|
||||
maxMessageBodySize + 64 * 1024 : Integer.MAX_VALUE;
|
||||
byteBuf = alloc.directBuffer(maxMessageSize);
|
||||
this.maxMessageBodySize = maxMessageBodySize;
|
||||
this.maxMessageSize = maxMessageSize;
|
||||
@@ -1692,6 +1704,18 @@ public class CommitLog {
|
||||
public ByteBuffer getEncoderBuffer() {
|
||||
return this.byteBuf.nioBuffer();
|
||||
}
|
||||
|
||||
public int getMaxMessageBodySize() {
|
||||
return this.maxMessageBodySize;
|
||||
}
|
||||
|
||||
public void updateEncoderBufferCapacity(int newMaxMessageBodySize) {
|
||||
this.maxMessageBodySize = newMaxMessageBodySize;
|
||||
//Reserve 64kb for encoding buffer outside body
|
||||
this.maxMessageSize = Integer.MAX_VALUE - newMaxMessageBodySize >= 64 * 1024 ?
|
||||
this.maxMessageBodySize + 64 * 1024 : Integer.MAX_VALUE;
|
||||
this.byteBuf.capacity(this.maxMessageSize);
|
||||
}
|
||||
}
|
||||
|
||||
static class PutMessageThreadLocal {
|
||||
|
||||
@@ -990,33 +990,17 @@ public class DefaultMessageStore implements MessageStore {
|
||||
long offset = queryOffsetResult.getPhyOffsets().get(m);
|
||||
|
||||
try {
|
||||
|
||||
boolean match = true;
|
||||
MessageExt msg = this.lookMessageByOffset(offset);
|
||||
if (0 == m) {
|
||||
lastQueryMsgTime = msg.getStoreTimestamp();
|
||||
}
|
||||
|
||||
// String[] keyArray = msg.getKeys().split(MessageConst.KEY_SEPARATOR);
|
||||
// if (topic.equals(msg.getTopic())) {
|
||||
// for (String k : keyArray) {
|
||||
// if (k.equals(key)) {
|
||||
// match = true;
|
||||
// break;
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
|
||||
if (match) {
|
||||
SelectMappedBufferResult result = this.commitLog.getData(offset, false);
|
||||
if (result != null) {
|
||||
int size = result.getByteBuffer().getInt(0);
|
||||
result.getByteBuffer().limit(size);
|
||||
result.setSize(size);
|
||||
queryMessageResult.addMessage(result);
|
||||
}
|
||||
} else {
|
||||
log.warn("queryMessage hash duplicate, {} {}", topic, key);
|
||||
SelectMappedBufferResult result = this.commitLog.getData(offset, false);
|
||||
if (result != null) {
|
||||
int size = result.getByteBuffer().getInt(0);
|
||||
result.getByteBuffer().limit(size);
|
||||
result.setSize(size);
|
||||
queryMessageResult.addMessage(result);
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.error("queryMessage exception", e);
|
||||
@@ -1387,12 +1371,8 @@ public class DefaultMessageStore implements MessageStore {
|
||||
private void checkSelf() {
|
||||
this.commitLog.checkSelf();
|
||||
|
||||
Iterator<Entry<String, ConcurrentMap<Integer, ConsumeQueue>>> it = this.consumeQueueTable.entrySet().iterator();
|
||||
while (it.hasNext()) {
|
||||
Entry<String, ConcurrentMap<Integer, ConsumeQueue>> next = it.next();
|
||||
Iterator<Entry<Integer, ConsumeQueue>> itNext = next.getValue().entrySet().iterator();
|
||||
while (itNext.hasNext()) {
|
||||
Entry<Integer, ConsumeQueue> cq = itNext.next();
|
||||
for (Entry<String, ConcurrentMap<Integer, ConsumeQueue>> next : this.consumeQueueTable.entrySet()) {
|
||||
for (Entry<Integer, ConsumeQueue> cq : next.getValue().entrySet()) {
|
||||
cq.getValue().checkSelf();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,8 +116,18 @@ public class MultiDispatch {
|
||||
}
|
||||
}
|
||||
|
||||
public void updateMaxMessageSize(CommitLog.PutMessageThreadLocal putMessageThreadLocal) {
|
||||
int newMaxMessageSize = this.messageStore.getMessageStoreConfig().getMaxMessageSize();
|
||||
if (newMaxMessageSize >= 10 &&
|
||||
putMessageThreadLocal.getEncoder().getMaxMessageBodySize() != newMaxMessageSize) {
|
||||
putMessageThreadLocal.getEncoder().updateEncoderBufferCapacity(newMaxMessageSize);
|
||||
}
|
||||
}
|
||||
|
||||
private boolean rebuildMsgInner(MessageExtBrokerInner msgInner) {
|
||||
MessageExtEncoder encoder = this.commitLog.getPutMessageThreadLocal().get().getEncoder();
|
||||
CommitLog.PutMessageThreadLocal putMessageThreadLocal = this.commitLog.getPutMessageThreadLocal().get();
|
||||
updateMaxMessageSize(putMessageThreadLocal);
|
||||
MessageExtEncoder encoder = putMessageThreadLocal.getEncoder();
|
||||
PutMessageResult encodeResult = encoder.encode(msgInner);
|
||||
if (encodeResult != null) {
|
||||
LOGGER.error("rebuild msgInner for multiDispatch", encodeResult);
|
||||
|
||||
@@ -696,6 +696,28 @@ public class DefaultMessageStoreTest {
|
||||
assertTrue(encodeResult5.getPutMessageStatus() == PutMessageStatus.MESSAGE_ILLEGAL);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDynamicMaxMessageSize(){
|
||||
MessageExtBrokerInner messageExtBrokerInner = buildMessage();
|
||||
MessageStoreConfig messageStoreConfig = ((DefaultMessageStore) messageStore).getMessageStoreConfig();
|
||||
int originMaxMessageSize = messageStoreConfig.getMaxMessageSize();
|
||||
|
||||
messageExtBrokerInner.setBody(new byte[originMaxMessageSize + 10]);
|
||||
PutMessageResult putMessageResult = messageStore.putMessage(messageExtBrokerInner);
|
||||
assertTrue(putMessageResult.getPutMessageStatus() == PutMessageStatus.MESSAGE_ILLEGAL);
|
||||
|
||||
int newMaxMessageSize = originMaxMessageSize + 10;
|
||||
messageStoreConfig.setMaxMessageSize(newMaxMessageSize);
|
||||
putMessageResult = messageStore.putMessage(messageExtBrokerInner);
|
||||
assertTrue(putMessageResult.getPutMessageStatus() == PutMessageStatus.PUT_OK);
|
||||
|
||||
messageStoreConfig.setMaxMessageSize(10);
|
||||
putMessageResult = messageStore.putMessage(messageExtBrokerInner);
|
||||
assertTrue(putMessageResult.getPutMessageStatus() == PutMessageStatus.MESSAGE_ILLEGAL);
|
||||
|
||||
messageStoreConfig.setMaxMessageSize(originMaxMessageSize);
|
||||
}
|
||||
|
||||
private class MyMessageArrivingListener implements MessageArrivingListener {
|
||||
@Override
|
||||
public void arriving(String topic, int queueId, long logicOffset, long tagsCode, long msgStoreTime,
|
||||
|
||||
@@ -382,5 +382,4 @@ public class DLedgerCommitlogTest extends MessageStoreTestBase {
|
||||
followerStore.shutdown();
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -46,8 +46,8 @@ public class DLedgerMultiPathTest extends MessageStoreTestBase {
|
||||
DefaultMessageStore dLedgerStore = createDLedgerMessageStore(base, group, "n0", peers, multiStorePath, null);
|
||||
Thread.sleep(2000);
|
||||
doPutMessages(dLedgerStore, topic, 0, 1000, 0);
|
||||
Assert.assertEquals(11, dLedgerStore.getMaxPhyOffset()/dLedgerStore.getMessageStoreConfig().getMappedFileSizeCommitLog());
|
||||
Thread.sleep(500);
|
||||
Assert.assertEquals(11, dLedgerStore.getMaxPhyOffset()/dLedgerStore.getMessageStoreConfig().getMappedFileSizeCommitLog());
|
||||
Assert.assertEquals(0, dLedgerStore.getMinOffsetInQueue(topic, 0));
|
||||
Assert.assertEquals(1000, dLedgerStore.getMaxOffsetInQueue(topic, 0));
|
||||
Assert.assertEquals(0, dLedgerStore.dispatchBehindBytes());
|
||||
|
||||
Reference in New Issue
Block a user