關(guān)于RabbitMQ的Channel默認(rèn)線程
前言
最近做了一個(gè)小功能,是通過一個(gè)客戶端消費(fèi)者監(jiān)聽隊(duì)列消息, 代碼如下:
Connection conn = getConnection(); Channel channel = conn.createChannel(); MessageConsumer consumer = ... channel.basicConsume(realQueue, true, consumer);
通過jvm工具觀察rabbitmq的線程使用情況,發(fā)現(xiàn)生產(chǎn)者每發(fā)一條消息,消費(fèi)者這邊就會(huì)創(chuàng)建一條線程, 言下之意,一個(gè)channel當(dāng)消息來到時(shí)就會(huì)異步處理這些消息.
定位
通過斷點(diǎn)查找發(fā)現(xiàn)原來是 ConsumerWorkService這個(gè)類控制的。
這個(gè)類顧名思義,就是消費(fèi)者工作 ExecutorService, 這里的Service表示的是ExecutorService
這個(gè)類構(gòu)造函數(shù)里有一個(gè)executor參數(shù),當(dāng)這個(gè)參數(shù)為空時(shí),就會(huì)創(chuàng)建一個(gè)Executors.newFixedThreadPool,代碼如下:
final public class ConsumerWorkService {
private static final int MAX_RUNNABLE_BLOCK_SIZE = 16;
private static final int DEFAULT_NUM_THREADS = Runtime.getRuntime().availableProcessors() * 2;
private final ExecutorService executor;
private final boolean privateExecutor;
private final WorkPool<Channel, Runnable> workPool;
private final int shutdownTimeout;
public ConsumerWorkService(ExecutorService executor, ThreadFactory threadFactory, int queueingTimeout, int shutdownTimeout) {
this.privateExecutor = (executor == null);
this.executor = (executor == null) ? Executors.newFixedThreadPool(DEFAULT_NUM_THREADS, threadFactory)
: executor;
this.workPool = new WorkPool<>(queueingTimeout);
this.shutdownTimeout = shutdownTimeout;
}
...默認(rèn)的executor 會(huì)使用 CPU核數(shù)的2倍 作為線程池里線程的數(shù)量。
所以到底是要用多個(gè)channel,還是單個(gè)channel,這個(gè)就是其中一個(gè)參考依據(jù)。
executor是怎么傳進(jìn)來的
答案:
ConnectionFactory -> AMQConnection -> ChannelManager -> ConsumerWorkService
ConnectionFactory有一個(gè)屬性是 shareExecutorService ,這個(gè)屬性表示內(nèi)部使用共享的唯一一個(gè)ExecutorService 設(shè)置這個(gè)屬性就可以一直傳到ConsumerWorkService中。
除了ConnectionFactory.setShareExecutorService方法以外, 還可以在Connection被創(chuàng)建時(shí),設(shè)置executorService ConnectionFactory的newConnection方法:
public Connection newConnection(ExecutorService executor) throws IOException, TimeoutException;
總結(jié)
通過設(shè)置shareExecutorService,無論多少個(gè)channel,都可以統(tǒng)一控制線程數(shù)量、隊(duì)列數(shù)量, 根據(jù)實(shí)際情況進(jìn)行配置。

public class RabbitMqUtil {
public static Channel getChannel() throws Exception{
//創(chuàng)建一個(gè)連接工廠
ConnectionFactory factory = new ConnectionFactory();
//連接服務(wù)器
factory.setHost("114.***.***.***");
//用戶名
factory.setUsername("admin");
//密碼
factory.setPassword("123");
//創(chuàng)建連接
// ExecutorService executor = Executors.newFixedThreadPool(1); 設(shè)置線程池中的個(gè)數(shù),把executor傳給newConnection()
Connection connection = factory.newConnection();
//獲取信道
Channel channel = connection.createChannel();
return channel;
}
}到此這篇關(guān)于關(guān)于RabbitMQ的Channel默認(rèn)線程的文章就介紹到這了,更多相關(guān)RabbitMQ的Channel內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
SpringCloud Eureka服務(wù)發(fā)現(xiàn)實(shí)現(xiàn)過程
這篇文章主要介紹了SpringCloud Eureka服務(wù)發(fā)現(xiàn)實(shí)現(xiàn)過程,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-11-11
Java中for(;;)和while(true)的區(qū)別
這篇文章主要介紹了 Java中for(;;)和while(true)的區(qū)別,文章圍繞for(;;)和while(true)的相關(guān)自來哦展開詳細(xì)內(nèi)容,需要的小伙伴可以參考一下,希望對(duì)大家有所幫助2021-11-11
Spring中@ExceptionHandler注解的使用方式
這篇文章主要介紹了Spring中@ExceptionHandler注解的使用方式,@ExceptionHandler注解我們一般是用來自定義異常的,可以認(rèn)為它是一個(gè)異常攔截器(處理器),需要的朋友可以參考下2024-01-01
詳解在Spring中如何自動(dòng)創(chuàng)建代理
這篇文章主要介紹了詳解在Spring中如何自動(dòng)創(chuàng)建代理,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧2018-07-07
SchedulingConfigurer實(shí)現(xiàn)動(dòng)態(tài)定時(shí),導(dǎo)致ApplicationRunner無效解決
這篇文章主要介紹了SchedulingConfigurer實(shí)現(xiàn)動(dòng)態(tài)定時(shí),導(dǎo)致ApplicationRunner無效的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-05-05
劍指Offer之Java算法習(xí)題精講字符串操作與數(shù)組及二叉搜索樹
跟著思路走,之后從簡單題入手,反復(fù)去看,做過之后可能會(huì)忘記,之后再做一次,記不住就反復(fù)做,反復(fù)尋求思路和規(guī)律,慢慢積累就會(huì)發(fā)現(xiàn)質(zhì)的變化2022-03-03
PowerDesigner連接數(shù)據(jù)庫的實(shí)例詳解
這篇文章主要介紹了PowerDesigner連接數(shù)據(jù)庫的實(shí)例詳解的相關(guān)資料,如有疑問請(qǐng)留言或者到本站社區(qū)交流討論,需要的朋友可以參考下2017-10-10

