RabbitMQ消費(fèi)端單線程與多線程案例講解
?? 一、基礎(chǔ)概念
| 模型 | 消費(fèi)者數(shù)量 | 每個(gè)消費(fèi)者內(nèi)部線程數(shù) | 順序性 | 場景說明 |
|---|---|---|---|---|
| 單消費(fèi)者單線程 | 1 | 1 | ? 保序 | 處理邏輯簡單,保證順序的常見場景 |
| 單消費(fèi)者多線程 | 1 | >1 | ? 不保序 | 提升處理能力,放棄順序要求 |
| 多消費(fèi)者單線程 | >1 | 1 | ? 不保序 | 多個(gè)隊(duì)列/分區(qū)消費(fèi),提升并發(fā) |
| 多消費(fèi)者多線程 | >1 | >1 | ? 不保序 | 高并發(fā)場景下批量處理,放棄順序 |
concurrency# 初始消費(fèi)者線程數(shù) max-concurrency# 最大消費(fèi)者線程數(shù) prefetch# 每個(gè)消費(fèi)者預(yù)取的消息數(shù)
concurrency: 2- 表示初始創(chuàng)建的消費(fèi)者線程數(shù)量
- 系統(tǒng)啟動(dòng)時(shí)會(huì)立即創(chuàng)建 2 個(gè)消費(fèi)者線程
- 這些線程會(huì)持續(xù)監(jiān)聽消息隊(duì)列
max-concurrency: 2- 表示允許的最大消費(fèi)者線程數(shù)量
- 這里設(shè)置為 2(與 concurrency 相同),表示線程數(shù)不會(huì)動(dòng)態(tài)擴(kuò)展
- 如果設(shè)置 max-concurrency > concurrency,系統(tǒng)會(huì)在負(fù)載高時(shí)動(dòng)態(tài)增加消費(fèi)者
詳細(xì)解釋:
concurrency和max-concurrency不會(huì)影響每個(gè)消費(fèi)者是否是多線程執(zhí)行,只會(huì)導(dǎo)致有多個(gè)消費(fèi)者線程,只有用線程池才會(huì)導(dǎo)致每個(gè)消費(fèi)者多線程消費(fèi)
而沒有用線程池,也設(shè)置prefetch是因?yàn)?/span>消息被大量預(yù)取,單線程處理不過來時(shí)堆積等待,單線程并不會(huì)影響消息的順序性,只有使用了線程池才會(huì)影響
使用了線程池一定會(huì)導(dǎo)致消息順序性問題這與設(shè)不設(shè)置prefetch無關(guān),因?yàn)槭褂镁€程池后,任務(wù)交個(gè)線程池就返回了屬于異步
舉個(gè)例子:
1. RabbitMQ 給消費(fèi)者推送消息1,消費(fèi)者收到,提交給線程池任務(wù)A(耗時(shí)長)。
2. 消費(fèi)者馬上ACK消息1(因?yàn)闃I(yè)務(wù)交給線程池了,自己處理完畢的感覺
3. RabbitMQ 再給消費(fèi)者推送消息2,消費(fèi)者收到,提交給線程池任務(wù)B(耗時(shí)短)。
4. R線程池調(diào)度先跑完任務(wù)B,后跑任務(wù)A。
? 單消費(fèi)者 + 單線程消費(fèi)
- 保證順序:消費(fèi)者內(nèi)部串行執(zhí)行。
- 配置關(guān)鍵:
spring:
rabbitmq:
listener:
simple:
concurrency: 1
max-concurrency: 1
prefetch: 1消費(fèi)者代碼:
@Component
public class MultiConsumerSingleThread {
@RabbitListener(queues = "order_queue", concurrency = "2")
public void receive(String message) {
System.out.println("?? [線程:" + Thread.currentThread().getName() + "] 收到消息:" + message);
try {
Thread.sleep(500);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}? 單消費(fèi)者 + 多線程消費(fèi)
- 不保順序:一個(gè)消費(fèi)者使用線程池異步處理消息。
- 配置關(guān)鍵:默認(rèn)配置 + 手動(dòng)異步處理
- 消費(fèi)者代碼:
@Component
public class MultiThreadConsumer {
private final ExecutorService executor = Executors.newFixedThreadPool(5);
@RabbitListener(queues = "order_queue")
public void receive(String message) {
executor.submit(() -> {
System.out.println("?? [線程:" + Thread.currentThread().getName() + "] 收到消息:" + message);
try {
Thread.sleep(500); // 模擬耗時(shí)
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}
}說明:消息提交到線程池,先到的不一定先處理完成,順序可能亂。
? 多消費(fèi)者 + 單線程消費(fèi)
- 不保順序:多個(gè)消費(fèi)者實(shí)例輪詢分配消息,各自順序保留,但整體順序錯(cuò)亂。
- 配置關(guān)鍵:
spring:
rabbitmq:
listener:
simple:
concurrency: 2
max-concurrency: 2
prefetch: 1消費(fèi)者代碼(共享類,也可拆成多個(gè)類模擬多實(shí)例):
@Component
public class MultiConsumerSingleThread {
//concurrency = "2":它和配置文件中的 concurrency: 2 作用一致,但優(yōu)先級(jí)更高。
@RabbitListener(queues = "order_queue", concurrency = "2")
public void receive(String message) {
System.out.println("?? [線程:" + Thread.currentThread().getName() + "] 收到消息:" + message);
try {
Thread.sleep(500);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}? 多消費(fèi)者 + 多線程消費(fèi)
- 不保順序:每個(gè)消費(fèi)者又使用線程池異步處理消息,最大吞吐量模式。
- 適合場景:數(shù)據(jù)導(dǎo)入、日志收集、發(fā)送通知等對順序無要求的批量處理。
- 配置關(guān)鍵:
spring:
rabbitmq:
listener:
simple:
concurrency: 3
max-concurrency: 3
prefetch: 10消費(fèi)者代碼:
@Component
public class MultiConsumerMultiThread {
private final ExecutorService executor = Executors.newFixedThreadPool(10);
@RabbitListener(queues = "order_queue", concurrency = "3")
public void receive(String message) {
executor.submit(() -> {
System.out.println("?? [線程:" + Thread.currentThread().getName() + "] 收到消息:" + message);
try {
Thread.sleep(500);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}
}?? 補(bǔ)充說明
concurrency: 控制并發(fā)消費(fèi)者數(shù)量,等于消費(fèi)者數(shù)。prefetch: 控制每個(gè)消費(fèi)者本地最多拉取多少條消息(如 1 表示嚴(yán)格串行處理)。- 每個(gè)
@RabbitListener本質(zhì)上是一個(gè)容器,可以通過concurrency配置“實(shí)例個(gè)數(shù)”。
到此這篇關(guān)于RabbitMQ消費(fèi)端單線程與多線程的文章就介紹到這了,更多相關(guān)RabbitMQ單線程與多線程內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
IDEA配置SpringBoot熱啟動(dòng),以及熱啟動(dòng)失效問題
這篇文章主要介紹了IDEA配置SpringBoot熱啟動(dòng),以及熱啟動(dòng)失效問題,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-11-11
Java中Stringbuilder和正則表達(dá)式示例詳解
Java語言為字符串連接運(yùn)算符(+)提供特殊支持,并為其他對象轉(zhuǎn)換為字符串,字符串連接是通過StringBuilder(或StringBuffer)類及其append方法實(shí)現(xiàn)的,這篇文章主要給大家介紹了關(guān)于Java中Stringbuilder和正則表達(dá)式的相關(guān)資料,需要的朋友可以參考下2024-02-02
如何基于Autowired對構(gòu)造函數(shù)進(jìn)行注釋
這篇文章主要介紹了如何基于Autowired對構(gòu)造函數(shù)進(jìn)行注釋,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-10-10
SpringBoot整合Mybatis自定義攔截器不起作用的處理方案
這篇文章主要介紹了SpringBoot整合Mybatis自定義攔截器不起作用的處理方案,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-09-09
Request的包裝類HttpServletRequestWrapper的使用說明
這篇文章主要介紹了Request的包裝類HttpServletRequestWrapper的使用說明,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-08-08

