不吃不喝看了好久才概括出這么一點點東西,希望大佬們能夠有耐心看一看,遇到說的不對的地方,也歡迎在評論區或者私信與我交流
另外完整版的代碼注釋,我在我的github上也添加了,感興趣的小伙伴也可以點擊這個鏈接去看一波 github地址
覺得我講的有那么一點點道理,對你有那么一丟丟的幫助的,也可以給我一波點贊關注666喲~

廢話不多說,下面開始我的表演~
RocketMQ全域流程圖

上來就是這么一大張圖片,相信大家肯定完全不想看下去,(那么我為什么還要放在一開始呢?主要是為了能夠讓大家有一個全域的印象,然后后續復習的時候也可以根據這個流程圖去具體復習)
那么,下面我們就針對一些問題來具體描述RocketMQ的作業流程 此處內容會不斷補充,也歡迎大家把遇到的問題在評論區留下來
訊息消費邏輯
訊息消費可以分為三大模塊
- Rebalance
- 拉取訊息
- 消費訊息
Rebalance

// RebalanceImpl
public void doRebalance(final boolean isOrder) {
Map<String, SubscriptionData> subTable = this.getSubscriptionInner();
if (subTable != null) {
// 遍歷每個主題的佇列
// subTable 會在 DefaultMQPushConsumerImpl 的 subscribe 和 unsubscribe 時修改
for (final Map.Entry<String, SubscriptionData> entry : subTable.entrySet()) {
final String topic = entry.getKey();
try {
// 對佇列進行重新負載
this.rebalanceByTopic(topic, isOrder);
} catch (Throwable e) {
if (!topic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
log.warn("rebalanceByTopic Exception", e);
}
}
}
}
this.truncateMessageQueueNotMyTopic();
}
private void rebalanceByTopic(final String topic, final boolean isOrder) {
switch (messageModel) {
case BROADCASTING: {
Set<MessageQueue> mqSet = this.topicSubscribeInfoTable.get(topic);
if (mqSet != null) {
boolean changed = this.updateProcessQueueTableInRebalance(topic, mqSet, isOrder);
if (changed) {
this.messageQueueChanged(topic, mqSet, mqSet);
log.info("messageQueueChanged {} {} {} {}",
consumerGroup,
topic,
mqSet,
mqSet);
}
} else {
log.warn("doRebalance, {}, but the topic[{}] not exist.", consumerGroup, topic);
}
break;
}
case CLUSTERING: {
// topicSubscribeInfoTable topic訂閱資訊快取表
Set<MessageQueue> mqSet = this.topicSubscribeInfoTable.get(topic);
// 發送請求到broker獲取topic下該消費組內當前所有的消費者客戶端id
List<String> cidAll = this.mQClientFactory.findConsumerIdList(topic, consumerGroup);
if (null == mqSet) {
if (!topic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
log.warn("doRebalance, {}, but the topic[{}] not exist.", consumerGroup, topic);
}
}
if (null == cidAll) {
log.warn("doRebalance, {} {}, get consumer id list failed", consumerGroup, topic);
}
if (mqSet != null && cidAll != null) {
List<MessageQueue> mqAll = new ArrayList<MessageQueue>();
mqAll.addAll(mqSet);
// 排序保證了同一個消費組內消費者看到的視圖保持一致,確保同一個消費佇列不會被多個消費者分配
Collections.sort(mqAll);
Collections.sort(cidAll);
// 分配演算法 (盡量使用前兩種)
// 默認有5種 1)平均分配 2)平均輪詢分配 3)一致性hash
// 4)根據配置 為每一個消費者配置固定的訊息佇列 5)根據broker部署機房名,對每個消費者負責不同的broker上的佇列
// 但是如果消費者數目大于訊息佇列數量,則會有些消費者無法消費訊息
AllocateMessageQueueStrategy strategy = this.allocateMessageQueueStrategy;
// 當前消費者分配到的佇列
List<MessageQueue> allocateResult = null;
try {
allocateResult = strategy.allocate(
this.consumerGroup,
this.mQClientFactory.getClientId(),
mqAll,
cidAll);
} catch (Throwable e) {
log.error("AllocateMessageQueueStrategy.allocate Exception. allocateMessageQueueStrategyName={}", strategy.getName(),
e);
return;
}
Set<MessageQueue> allocateResultSet = new HashSet<MessageQueue>();
if (allocateResult != null) {
allocateResultSet.addAll(allocateResult);
}
// 更新訊息消費佇列,如果是新增的訊息消費佇列,則會創建一個訊息拉取請求并立即執行拉取
boolean changed = this.updateProcessQueueTableInRebalance(topic, allocateResultSet, isOrder);
if (changed) {
log.info(
"rebalanced result changed. allocateMessageQueueStrategyName={}, group={}, topic={}, clientId={}, mqAllSize={}, cidAllSize={}, rebalanceResultSize={}, rebalanceResultSet={}",
strategy.getName(), consumerGroup, topic, this.mQClientFactory.getClientId(), mqSet.size(), cidAll.size(),
allocateResultSet.size(), allocateResultSet);
this.messageQueueChanged(topic, mqSet, allocateResultSet);
}
}
break;
}
default:
break;
}
}
private boolean updateProcessQueueTableInRebalance(final String topic, final Set<MessageQueue> mqSet,
final boolean isOrder) {
boolean changed = false;
Iterator<Entry<MessageQueue, ProcessQueue>> it = this.processQueueTable.entrySet().iterator();
while (it.hasNext()) {
Entry<MessageQueue, ProcessQueue> next = it.next();
MessageQueue mq = next.getKey();
ProcessQueue pq = next.getValue();
if (mq.getTopic().equals(topic)) {
// 當前分配到的佇列中不包含原先的佇列(說明當前佇列被分配給了其他消費者)
if (!mqSet.contains(mq)) {
// 丟棄 processQueue
pq.setDropped(true);
// 移除當前訊息佇列
if (this.removeUnnecessaryMessageQueue(mq, pq)) {
it.remove();
changed = true;
log.info("doRebalance, {}, remove unnecessary mq, {}", consumerGroup, mq);
}
} else if (pq.isPullExpired()) {
switch (this.consumeType()) {
case CONSUME_ACTIVELY:
break;
case CONSUME_PASSIVELY:
pq.setDropped(true);
if (this.removeUnnecessaryMessageQueue(mq, pq)) {
it.remove();
changed = true;
log.error("[BUG]doRebalance, {}, remove unnecessary mq, {}, because pull is pause, so try to fixed it",
consumerGroup, mq);
}
break;
default:
break;
}
}
}
}
List<PullRequest> pullRequestList = new ArrayList<PullRequest>();
for (MessageQueue mq : mqSet) {
// 訊息消費佇列快取中不存在當前佇列 本次分配新增的佇列
if (!this.processQueueTable.containsKey(mq)) {
// 向broker發起鎖定佇列請求 (向broker端請求鎖定MessageQueue,同時在本地鎖定對應的ProcessQueue)
if (isOrder && !this.lock(mq)) {
log.warn("doRebalance, {}, add a new mq failed, {}, because lock failed", consumerGroup, mq);
// 加鎖失敗,跳過,等待下一次佇列重新負載時再嘗試加鎖
continue;
}
// 從記憶體中移除該訊息佇列的消費進度
this.removeDirtyOffset(mq);
ProcessQueue pq = new ProcessQueue();
long nextOffset = -1L;
try {
nextOffset = this.computePullFromWhereWithException(mq);
} catch (Exception e) {
log.info("doRebalance, {}, compute offset failed, {}", consumerGroup, mq);
continue;
}
if (nextOffset >= 0) {
ProcessQueue pre = this.processQueueTable.putIfAbsent(mq, pq);
if (pre != null) {
log.info("doRebalance, {}, mq already exists, {}", consumerGroup, mq);
} else {
// 首次添加,構建拉取訊息的請求
log.info("doRebalance, {}, add a new mq, {}", consumerGroup, mq);
PullRequest pullRequest = new PullRequest();
pullRequest.setConsumerGroup(consumerGroup);
pullRequest.setNextOffset(nextOffset);
pullRequest.setMessageQueue(mq);
pullRequest.setProcessQueue(pq);
pullRequestList.add(pullRequest);
changed = true;
}
} else {
log.warn("doRebalance, {}, add new mq failed, {}", consumerGroup, mq);
}
}
}
// 立即拉取訊息(對新增的佇列)
this.dispatchPullRequest(pullRequestList);
return changed;
}
由流程圖和代碼,我們可以得知,集群模式下訊息負載主要有以下幾個步驟:
- 從Broker獲取訂閱當前Topic的消費者串列
- 根據具體的策略進行負載均衡
- 對當前消費者分配到的佇列進行處理
- 原來有,現在沒有:丟棄對應的訊息處理佇列(ProcessQueue)
- 原來沒有,現在有:添加訊息處理佇列(ProcessQueue),如果是第一次新增,還會創建一個訊息拉取請求
拉取訊息

拉取訊息的代碼太多了,我就不再這里貼出來了,
我在這里說一下大致流程,然后有幾個需要注意的地方
流程:在我們Rebalance第一次添加負責的佇列和后續拉取訊息后,都會再提交一個拉取請求到拉取請求佇列(pullRequestQueue)中,然后有一個執行緒不停的去里面獲取拉取請求,去執行拉取的操作
這里說一個RocketMQ消費者這邊設計的一個亮點
它將拉取訊息,消費訊息通過兩個任務佇列的方式進行解耦,然后每一個模塊僅需要負責它自己的功能,(雖然大佬們覺得很常見,但是當時我看的時候還是感覺妙呀~)
另外還有一點需要注意的是:拉取訊息的時候broker和consumer都會對訊息進行過濾,只不過broker是根據tag的hash進行過濾的,而consumer是根據具體的tag字串匹配過濾的,這也是有的時候,明明拉取到了訊息,但是卻沒有需要消費的訊息產生的原因
既然說到了訊息過濾,這邊先簡單提一下RocketMQ訊息過濾的幾種方式
- 運算式過濾
- tag
- SQL92
- 類過濾
消費訊息

這邊也先說幾個注意點吧,后面再單獨出篇文章,
(一)順序消費和非順序消費消費失敗的處理
(二)消費失敗偏移量的更新:只有當前這批訊息全部消費成功后,才會將偏移量更新成為這批訊息最后一條的偏移量
(三)廣播訊息失敗不會重試,僅列印失敗日志
補充:為什么同一個消費組下消費者的訂閱資訊要相同
首先,先說一下什么叫做同一個消費組下消費者的訂閱資訊要相同
即:在相同的GroupId下,每一個消費者他們的訂閱內容(Topic+Tag)要保持一致,否則會導致訊息無法被正常消費
參考檔案:阿里云:訂閱關系一致

我們在看待這個問題的時候,可以把它分為兩類情況考慮
- topic不一致
- tag不一致
(一)topic不一致的問題
首先先說一個場景,消費者A監聽了TopicA,消費者B監聽了TopicB,但是消費者A和消費者B同屬一個groupTest
在Rebalance階段,消費者A對TopicA進行負載均衡時,會去查詢groupTest下的所有消費者資訊,獲取到了消費者A和消費者B,此時就會將TopicA的佇列對消費者A和消費者B進行負載均衡(例如消費者A分配到了1234四個佇列,消費者B分配到了5678四個佇列),此時消費者B沒有針對TopicA的處理邏輯,就會導致推送到5678這幾個佇列里面的訊息沒有辦法得到處理,
(二)tag不一致的問題
隨著消費者A,消費者B負載均衡的不斷進行,會不斷把最新的訂閱資訊(訊息過濾規則)上報給broker,broker就會不斷的覆寫更新,導致tag資訊不停地變化,而tag的變化在消費者拉取訊息時broker的過濾就會產生影響,會導致一些本來要被消費者拉取到的訊息被broker過濾掉
延時佇列是如何作業的

由流程圖中我們不難看出,RocketMQ對延時訊息的處理,是交由Timer去完成的(相關類ScheduleMessageService),在Timer的任務佇列中讀取需要處理的延遲任務,將訊息從延遲佇列轉發到具體的業務佇列中
此處補充一點:此處提到的Timer為java工具類包(java.util.Timer)下的一個定時任務工具,它主要由兩個部分:TaskQueue queue(任務佇列)和TimerThread thread(作業執行緒),這邊我把它簡單的類比為一個單執行緒的作業執行緒池
另外在ScheduleMessageService中使用到了Timer的兩個方法,我在這里先單獨列出來下
- this.timer.schedule :在任務執行成功后,再加上對應的周期,然后再執行
- this.timer.scheduleAtFixedRate :每隔指定時間就執行一次,與任務執行時間無關
話不多少,貼上原始碼*(原始碼雖然枯燥,但希望可以耐心的看完)*
// ScheduleMessageService
public void start() {
if (started.compareAndSet(false, true)) {
super.load();
this.timer = new Timer("ScheduleMessageTimerThread", true);
// 根據延時佇列創建對應的定時任務
for (Map.Entry<Integer, Long> entry : this.delayLevelTable.entrySet()) {
Integer level = entry.getKey();
Long timeDelay = entry.getValue();
Long offset = this.offsetTable.get(level);
if (null == offset) {
offset = 0L;
}
if (timeDelay != null) {
// 第一次,延遲一秒執行任務,后續根據對應延時時間來執行
// 延時級別和訊息佇列id對應關系 : 訊息佇列id = 延時級別 - 1
// shedule 在任務執行成功后,再加上對應的周期,然后再執行
this.timer.schedule(new DeliverDelayedMessageTimerTask(level, offset), FIRST_DELAY_TIME);
}
}
// scheduleAtFixedRate 每隔指定時間就執行一次,與任務執行時間無關
this.timer.scheduleAtFixedRate(new TimerTask() {
@Override
public void run() {
try {
if (started.get()) {
// 每個十秒持久化一次延遲佇列的處理進度
ScheduleMessageService.this.persist();
}
} catch (Throwable e) {
log.error("scheduleAtFixedRate flush exception", e);
}
}
}, 10000, this.defaultMessageStore.getMessageStoreConfig().getFlushDelayOffsetInterval());
}
}
// DeliverDelayedMessageTimerTask
@Override
public void run() {
try {
if (isStarted()) {
this.executeOnTimeup();
}
} catch (Exception e) {
// XXX: warn and notify me
log.error("ScheduleMessageService, executeOnTimeup exception", e);
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(
this.delayLevel, this.offset), DELAY_FOR_A_PERIOD);
}
}
public void executeOnTimeup() {
// 根據 延時佇列topic 和 延時佇列id 查找消費佇列
ConsumeQueue cq =
ScheduleMessageService.this.defaultMessageStore.findConsumeQueue(TopicValidator.RMQ_SYS_SCHEDULE_TOPIC,
delayLevel2QueueId(delayLevel));
long failScheduleOffset = offset;
if (cq != null) {
SelectMappedBufferResult bufferCQ = cq.getIndexBuffer(this.offset);
if (bufferCQ != null) {
try {
long nextOffset = offset;
int i = 0;
// 遍歷ConsumeQueue,每一個標準的ConsumeQueue條目為20位元組
ConsumeQueueExt.CqExtUnit cqExtUnit = new ConsumeQueueExt.CqExtUnit();
for (; i < bufferCQ.getSize(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) {
long offsetPy = bufferCQ.getByteBuffer().getLong();
int sizePy = bufferCQ.getByteBuffer().getInt();
long tagsCode = bufferCQ.getByteBuffer().getLong();
if (cq.isExtAddr(tagsCode)) {
if (cq.getExt(tagsCode, cqExtUnit)) {
tagsCode = cqExtUnit.getTagsCode();
} else {
//can't find ext content.So re compute tags code.
log.error("[BUG] can't find consume queue extend file content!addr={}, offsetPy={}, sizePy={}",
tagsCode, offsetPy, sizePy);
long msgStoreTime = defaultMessageStore.getCommitLog().pickupStoreTimestamp(offsetPy, sizePy);
tagsCode = computeDeliverTimestamp(delayLevel, msgStoreTime);
}
}
long now = System.currentTimeMillis();
long deliverTimestamp = this.correctDeliverTimestamp(now, tagsCode);
nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
// > 0 未到訊息消費時間
long countdown = deliverTimestamp - now;
if (countdown <= 0) {
MessageExt msgExt =
ScheduleMessageService.this.defaultMessageStore.lookMessageByOffset(
offsetPy, sizePy);
if (msgExt != null) {
try {
MessageExtBrokerInner msgInner = this.messageTimeup(msgExt);
if (TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC.equals(msgInner.getTopic())) {
log.error("[BUG] the real topic of schedule msg is {}, discard the msg. msg={}",
msgInner.getTopic(), msgInner);
continue;
}
// 放到對應的 %RETRY%+gid 重試topic下進行消費(轉發訊息)
PutMessageResult putMessageResult =
ScheduleMessageService.this.writeMessageStore
.putMessage(msgInner);
if (putMessageResult != null
&& putMessageResult.getPutMessageStatus() == PutMessageStatus.PUT_OK) {
if (ScheduleMessageService.this.defaultMessageStore.getMessageStoreConfig().isEnableScheduleMessageStats()) {
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incQueueGetNums(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, delayLevel - 1, putMessageResult.getAppendMessageResult().getMsgNum());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incQueueGetSize(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, delayLevel - 1, putMessageResult.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incGroupGetNums(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, putMessageResult.getAppendMessageResult().getMsgNum());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incGroupGetSize(MixAll.SCHEDULE_CONSUMER_GROUP, TopicValidator.RMQ_SYS_SCHEDULE_TOPIC, putMessageResult.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incTopicPutNums(msgInner.getTopic(), putMessageResult.getAppendMessageResult().getMsgNum(), 1);
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incTopicPutSize(msgInner.getTopic(),
putMessageResult.getAppendMessageResult().getWroteBytes());
ScheduleMessageService.this.defaultMessageStore.getBrokerStatsManager().incBrokerPutNums(putMessageResult.getAppendMessageResult().getMsgNum());
}
continue;
} else {
// XXX: warn and notify me
log.error(
"ScheduleMessageService, a message time up, but reput it failed, topic: {} msgId {}",
msgExt.getTopic(), msgExt.getMsgId());
ScheduleMessageService.this.timer.schedule(
new DeliverDelayedMessageTimerTask(this.delayLevel,
nextOffset), DELAY_FOR_A_PERIOD);
ScheduleMessageService.this.updateOffset(this.delayLevel,
nextOffset);
return;
}
} catch (Exception e) {
/*
* XXX: warn and notify me
*/
log.error(
"ScheduleMessageService, messageTimeup execute error, drop it. msgExt={}, nextOffset={}, offsetPy={}, sizePy={}", msgExt, nextOffset, offsetPy, sizePy, e);
}
}
} else {
// 會將下次任務執行時間設定為countdown 即 訊息的延時轉發時間-當前時間
ScheduleMessageService.this.timer.schedule(
new DeliverDelayedMessageTimerTask(this.delayLevel, nextOffset),
countdown);
ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset);
return;
}
} // end of for
// 更新延時佇列拉取任務進度
nextOffset = offset + (i / ConsumeQueue.CQ_STORE_UNIT_SIZE);
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(
this.delayLevel, nextOffset), DELAY_FOR_A_WHILE);
ScheduleMessageService.this.updateOffset(this.delayLevel, nextOffset);
return;
} finally {
bufferCQ.release();
}
} // end of if (bufferCQ != null)
else {
// 消費佇列不存在,默認為沒有需要消費的任務,跳過本次消費
long cqMinOffset = cq.getMinOffsetInQueue();
long cqMaxOffset = cq.getMaxOffsetInQueue();
if (offset < cqMinOffset) {
// 下次拉取任務進度更新
failScheduleOffset = cqMinOffset;
log.error("schedule CQ offset invalid. offset={}, cqMinOffset={}, cqMaxOffset={}, queueId={}",
offset, cqMinOffset, cqMaxOffset, cq.getQueueId());
}
if (offset > cqMaxOffset) {
failScheduleOffset = cqMaxOffset;
log.error("schedule CQ offset invalid. offset={}, cqMinOffset={}, cqMaxOffset={}, queueId={}",
offset, cqMinOffset, cqMaxOffset, cq.getQueueId());
}
}
} // end of if (cq != null)
// 根據延時等級創建一個任務
ScheduleMessageService.this.timer.schedule(new DeliverDelayedMessageTimerTask(this.delayLevel,
failScheduleOffset), DELAY_FOR_A_WHILE);
}
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/376998.html
標籤:其他
下一篇:windows 中 elasticsearch 7.15.2 和 插件 詳細 下載和安裝(解決elasticsearch.exceptions.ConnectionError)
