最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RabbitMQ消費(fèi)端單線程與多線程案例講解

 更新時(shí)間:2025年07月26日 09:30:25   作者:你我約定有三  
文章解析RabbitMQ消費(fèi)端單線程與多線程處理機(jī)制,說明concurrency控制消費(fèi)者數(shù)量,max-concurrency控制最大線程數(shù),prefetch影響消息預(yù)取量,強(qiáng)調(diào)線程池使用會(huì)導(dǎo)致順序混亂,適用于對順序無要求的批量處理場景,感興趣的朋友一起看看吧

?? 一、基礎(chǔ)概念

模型消費(fèi)者數(shù)量每個(gè)消費(fèi)者內(nèi)部線程數(shù)順序性場景說明
單消費(fèi)者單線程11? 保序處理邏輯簡單,保證順序的常見場景
單消費(fèi)者多線程1>1? 不保序提升處理能力,放棄順序要求
多消費(fèi)者單線程>11? 不保序多個(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)文章

  • Java實(shí)現(xiàn)全圖背景水印的示例詳解

    Java實(shí)現(xiàn)全圖背景水印的示例詳解

    這篇文章主要為大家詳細(xì)介紹了如何利用Java實(shí)現(xiàn)全圖背景水印的方法,文中的示例代碼講解詳細(xì),具有一定的借鑒價(jià)值,需要的可以參考一下
    2023-02-02
  • IDEA配置SpringBoot熱啟動(dòng),以及熱啟動(dòng)失效問題

    IDEA配置SpringBoot熱啟動(dòng),以及熱啟動(dòng)失效問題

    這篇文章主要介紹了IDEA配置SpringBoot熱啟動(dòng),以及熱啟動(dòng)失效問題,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-11-11
  • Java中Stringbuilder和正則表達(dá)式示例詳解

    Java中Stringbuilder和正則表達(dá)式示例詳解

    Java語言為字符串連接運(yùn)算符(+)提供特殊支持,并為其他對象轉(zhuǎn)換為字符串,字符串連接是通過StringBuilder(或StringBuffer)類及其append方法實(shí)現(xiàn)的,這篇文章主要給大家介紹了關(guān)于Java中Stringbuilder和正則表達(dá)式的相關(guān)資料,需要的朋友可以參考下
    2024-02-02
  • MyBatis中#{}和${}的區(qū)別詳解

    MyBatis中#{}和${}的區(qū)別詳解

    mybatis和ibatis總體來講都差不多的。下面小編給大家探討下mybatis中#{}和${}的區(qū)別,感興趣的朋友一起學(xué)習(xí)吧
    2016-08-08
  • Java編程中的HashSet和BitSet詳解

    Java編程中的HashSet和BitSet詳解

    這篇文章主要介紹了Java編程中的HashSet和BitSet詳解的相關(guān)資料,需要的朋友可以參考下
    2017-03-03
  • Java List移除相應(yīng)元素的超簡潔寫法分享

    Java List移除相應(yīng)元素的超簡潔寫法分享

    這篇文章主要介紹了Java List移除相應(yīng)元素的超簡潔寫法,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • 如何基于Autowired對構(gòu)造函數(shù)進(jìn)行注釋

    如何基于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自定義攔截器不起作用的處理方案

    這篇文章主要介紹了SpringBoot整合Mybatis自定義攔截器不起作用的處理方案,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Request的包裝類HttpServletRequestWrapper的使用說明

    Request的包裝類HttpServletRequestWrapper的使用說明

    這篇文章主要介紹了Request的包裝類HttpServletRequestWrapper的使用說明,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • 詳解Java 類的加載機(jī)制

    詳解Java 類的加載機(jī)制

    這篇文章主要介紹了Java 類的加載機(jī)制,幫助大家更好的理解和學(xué)習(xí)Java,感興趣的朋友可以了解下
    2020-08-08

最新評論

陵水| 安义县| 迁安市| 永兴县| 闸北区| 丰都县| 山丹县| 滦南县| 民乐县| 彭阳县| 乌兰察布市| 秦皇岛市| 潮安县| 五台县| 绥江县| 资兴市| 黎城县| 遵义县| 无棣县| 岳阳县| 永登县| 舒兰市| 绥江县| 万安县| 秦安县| 灵璧县| 海原县| 逊克县| 汕头市| 义乌市| 靖江市| 田阳县| 门头沟区| 磐石市| 卢氏县| 巴青县| 那曲县| 阆中市| 托克逊县| 永川市| 吉隆县|