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

docker啟動rabbitmq以及使用方式詳解

 更新時間:2022年08月04日 11:34:00   作者:Maackia  
RabbitMQ是一個由erlang開發(fā)的消息隊列,下面這篇文章主要給大家介紹了關(guān)于docker啟動rabbitmq以及使用的相關(guān)資料,文中通過圖文介紹的非常詳細,需要的朋友可以參考下

搜索rabbitmq鏡像

docker search rabbitmq:management

下載鏡像

docker pull rabbitmq:management

啟動容器

docker run -d --hostname localhost --name rabbitmq -p 15672:15672 -p 5672:5672 rabbitmq:management

打印容器

docker logs rabbitmq

訪問RabbitMQ Management

http://localhost:15672

賬戶密碼默認:guest

編寫生產(chǎn)者類

package com.xun.rabbitmqdemo.example;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;

public class Producer {
    private final static String QUEUE_NAME = "hello";
    public static void main(String[] args) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setUsername("guest");
        factory.setPassword("guest");
        factory.setHost("localhost");
        factory.setPort(5672);
        factory.setVirtualHost("/");

        Connection connection = factory.newConnection();

        Channel channel = connection.createChannel();
        /**
         * 生成一個queue隊列
         * 1、隊列名稱 QUEUE_NAME
         * 2、隊列里面的消息是否持久化(默認消息存儲在內(nèi)存中)
         * 3、該隊列是否只供一個Consumer消費 是否共享 設(shè)置為true可以多個消費者消費
         * 4、是否自動刪除 最后一個消費者斷開連接后 該隊列是否自動刪除
         * 5、其他參數(shù)
         */
        channel.queueDeclare(QUEUE_NAME,false,false,false,null);
        String message = "Hello world!";
        /**
         * 發(fā)送一個消息
         * 1、發(fā)送到哪個exchange交換機
         * 2、路由的key
         * 3、其他的參數(shù)信息
         * 4、消息體
         */
        channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
        System.out.println(" [x] Sent '"+message+"'");

        channel.close();
        connection.close();
    }
}

運行該方法,可以看到控制臺的打印

name=hello的隊列收到Message

消費者

package com.xun.rabbitmqdemo.example;

import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;

public class Receiver {
    private final static String QUEUE_NAME = "hello";
    public static void main(String[] args) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setUsername("guest");
        factory.setPassword("guest");
        factory.setHost("localhost");
        factory.setPort(5672);
        factory.setVirtualHost("/");
        factory.setConnectionTimeout(600000);//milliseconds
        factory.setRequestedHeartbeat(60);//seconds
        factory.setHandshakeTimeout(6000);//milliseconds
        factory.setRequestedChannelMax(5);
        factory.setNetworkRecoveryInterval(500);

        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        channel.queueDeclare(QUEUE_NAME,false,false,false,null);
        System.out.println("Waiting for messages. ");

        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties,byte[] body) throws IOException {
                String message = new String(body, "UTF-8");
                System.out.println(" [x] Received '" + message + "'");
            }
        };
        channel.basicConsume(QUEUE_NAME,true,consumer);
    }
}

工作隊列

RabbitMqUtils工具類

package com.xun.rabbitmqdemo.utils;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class RabbitMqUtils {
    public static Channel getChannel() throws Exception{
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        factory.setUsername("guest");
        factory.setPassword("guest");
        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();
        return channel;
    }
}

啟動2個工作線程

package com.xun.rabbitmqdemo.workQueue;

import com.rabbitmq.client.CancelCallback;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;
import com.xun.rabbitmqdemo.utils.RabbitMqUtils;

public class Work01 {
    private static final String QUEUE_NAME = "hello";
    public static void main(String[] args) throws Exception{
        Channel channel = RabbitMqUtils.getChannel();
        DeliverCallback deliverCallback = (consumerTag,delivery)->{
            String receivedMessage = new String(delivery.getBody());
            System.out.println("接收消息:"+receivedMessage);
        };
        CancelCallback cancelCallback = (consumerTag)->{
            System.out.println(consumerTag+"消費者取消消費接口回調(diào)邏輯");
        };
        System.out.println("C1 消費者啟動等待消費....");
        /**
         * 消費者消費消息
         * 1、消費哪個隊列
         * 2、消費成功后是否自動應答
         * 3、消費的接口回調(diào)
         * 4、消費未成功的接口回調(diào)
         */
        channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
    }
}
package com.xun.rabbitmqdemo.workQueue;

import com.rabbitmq.client.CancelCallback;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;
import com.xun.rabbitmqdemo.utils.RabbitMqUtils;

public class Work02 {
    private static final String QUEUE_NAME = "hello";
    public static void main(String[] args) throws Exception{
        Channel channel = RabbitMqUtils.getChannel();
        DeliverCallback deliverCallback = (consumerTag,delivery)->{
            String receivedMessage = new String(delivery.getBody());
            System.out.println("接收消息:"+receivedMessage);
        };
        CancelCallback cancelCallback = (consumerTag)->{
            System.out.println(consumerTag+"消費者取消消費接口回調(diào)邏輯");
        };
        System.out.println("C2 消費者啟動等待消費....");
        channel.basicConsume(QUEUE_NAME,true,deliverCallback,cancelCallback);
    }
}

啟動工作線程

啟動發(fā)送線程

package com.xun.rabbitmqdemo.workQueue;

import com.rabbitmq.client.Channel;
import com.xun.rabbitmqdemo.utils.RabbitMqUtils;
import java.util.Scanner;

public class Task01 {
    private static final String QUEUE_NAME = "hello";
    public static void main(String[] args) throws Exception{
        try(Channel channel= RabbitMqUtils.getChannel();){
            channel.queueDeclare(QUEUE_NAME,false,false,false,null);
            //從控制臺接收消息
            Scanner scanner = new Scanner(System.in);
            while(scanner.hasNext()){
                String message = scanner.next();
                channel.basicPublish("",QUEUE_NAME,null,message.getBytes());
                System.out.println("發(fā)送消息完成:"+message);
            }
        }
    }
}

啟動發(fā)送線程,此時發(fā)送線程等待鍵盤輸入

發(fā)送4個消息

可以看到2個工作線程按照順序分別接收message。

消息應答機制

rabbitmq將message發(fā)送給消費者后,就會將該消息標記為刪除。

但消費者在處理message過程中宕機,會導致消息的丟失。

因此需要設(shè)置手動應答。

生產(chǎn)者

import com.xun.rabbitmqdemo.utils.RabbitMqUtils;
import java.util.Scanner;

public class Task02 {
    private static final String TASK_QUEUE_NAME = "ack_queue";
    public static void main(String[] args) throws Exception{
        try(Channel channel = RabbitMqUtils.getChannel()){
            channel.queueDeclare(TASK_QUEUE_NAME,false,false,false,null);
            Scanner scanner = new Scanner(System.in);
            System.out.println("請輸入信息");
            while(scanner.hasNext()){
                String message = scanner.nextLine();
                channel.basicPublish("",TASK_QUEUE_NAME,null,message.getBytes());
                System.out.println("生產(chǎn)者task02發(fā)出消息"+ message);
            }
        }
    }
}

消費者

package com.xun.rabbitmqdemo.workQueue;
import com.rabbitmq.client.CancelCallback;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;
import com.xun.rabbitmqdemo.utils.RabbitMqUtils;
import com.xun.rabbitmqdemo.utils.SleepUtils;

public class Work03 {
    private static final String ACK_QUEUE_NAME = "ack_queue";
    public static void main(String[] args) throws Exception{
        Channel channel = RabbitMqUtils.getChannel();
        System.out.println("Work03 等待接收消息處理時間較短");
        DeliverCallback deliverCallback = (consumerTag,delivery)->{
            String message = new String(delivery.getBody());
            SleepUtils.sleep(1);
            System.out.println("接收到消息:"+message);
            /**
             * 1、消息的標記tag
             * 2、是否批量應答
             */
            channel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);
        };
        CancelCallback cancelCallback = (consumerTag)->{
            System.out.println(consumerTag+"消費者取消消費接口回調(diào)邏輯");
        };
        //采用手動應答
        boolean autoAck = false;
        channel.basicConsume(ACK_QUEUE_NAME,autoAck,deliverCallback,cancelCallback);
    }
}
package com.xun.rabbitmqdemo.workQueue;
import com.rabbitmq.client.CancelCallback;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;
import com.xun.rabbitmqdemo.utils.RabbitMqUtils;
import com.xun.rabbitmqdemo.utils.SleepUtils;

public class Work04 {
    private static final String ACK_QUEUE_NAME = "ack_queue";
    public static void main(String[] args) throws Exception{
        Channel channel = RabbitMqUtils.getChannel();
        System.out.println("Work04 等待接收消息處理時間較長");
        DeliverCallback deliverCallback = (consumerTag,delivery)->{
            String message = new String(delivery.getBody());
            SleepUtils.sleep(30);
            System.out.println("接收到消息:"+message);
            channel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);
        };
        CancelCallback cancelCallback = (consumerTag)->{
            System.out.println(consumerTag+"消費者取消消費接口回調(diào)邏輯");
        };
        //采用手動應答
        boolean autoAck = false;
        channel.basicConsume(ACK_QUEUE_NAME,autoAck,deliverCallback,cancelCallback);
    }
}

工具類SleepUtils

package com.xun.rabbitmqdemo.utils;
public class SleepUtils {
    public static void sleep(int second){
        try{
            Thread.sleep(1000*second);
        }catch (InterruptedException _ignored){
            Thread.currentThread().interrupt();
        }
    }
}

模擬

work04等待30s后發(fā)出ack

在work04處理message時手動停止線程,可以看到message:dd被rabbitmq交給了work03

不公平分發(fā)

上面的輪詢分發(fā),生產(chǎn)者依次向消費者按順序發(fā)送消息,但當消費者A處理速度很快,而消費者B處理速度很慢時,這種分發(fā)策略顯然是不合理的。
不公平分發(fā):

int prefetchCount = 1;
channel.basicQos(prefetchCount);

通過此配置,當消費者未處理完當前消息,rabbitmq會優(yōu)先將該message分發(fā)給空閑消費者。

總結(jié) 

到此這篇關(guān)于docker啟動rabbitmq以及使用的文章就介紹到這了,更多相關(guān)docker啟動rabbitmq及使用內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Docker?制作tomcat鏡像并部署項目的步驟

    Docker?制作tomcat鏡像并部署項目的步驟

    這篇文章主要介紹了Docker?制作tomcat鏡像并部署項目?,講解如何制作自己的tomcat鏡像,并使用tomcat部署項目,需要的朋友可以參考下
    2022-10-10
  • docker?部署?時序數(shù)據(jù)庫TDengine的思路詳解

    docker?部署?時序數(shù)據(jù)庫TDengine的思路詳解

    TDengineGUI是一個基于electron構(gòu)建的,針對時序數(shù)據(jù)庫TDengine的圖形化管理工具,這篇文章主要介紹了docker?部署?時序數(shù)據(jù)庫TDengine的思路詳解,需要的朋友可以參考下
    2025-04-04
  • 遷移docker鏡像到新服務(wù)器的具體操作流程

    遷移docker鏡像到新服務(wù)器的具體操作流程

    在日常工作中,我們有時會需要將服務(wù)器A上的鏡像上傳至服務(wù)器B上,這篇文章主要介紹了遷移docker鏡像到新服務(wù)器的具體操作流程,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2025-07-07
  • Docker?安裝Nginx與配置Nginx的案例

    Docker?安裝Nginx與配置Nginx的案例

    Nginx是一個高性能的HTTP和反向代理web服務(wù)器,ginx是一款輕量級的Web?服務(wù)器/反向代理服務(wù)器及電子郵件(IMAP/POP3)代理服務(wù)器,在BSD-like?協(xié)議下發(fā)行,下面通過本文給大家介紹Docker?安裝Nginx與配置Nginx的案例,感興趣的朋友一起看看吧
    2024-08-08
  • Docker搭建Jenkins并自動化打包部署項目的步驟

    Docker搭建Jenkins并自動化打包部署項目的步驟

    本文主要介紹了Docker搭建Jenkins并自動化打包部署項目的步驟,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-03-03
  • Docker實現(xiàn)將鏡像從1.2GB壓縮到200MB的優(yōu)化指南

    Docker實現(xiàn)將鏡像從1.2GB壓縮到200MB的優(yōu)化指南

    作為一名在容器化領(lǐng)域摸爬滾打多年的開發(fā)者,深知Docker鏡像大小對生產(chǎn)環(huán)境的影響,本文將詳細記錄優(yōu)化歷程,從問題分析到解決方案實施,從理論原理到實戰(zhàn)技巧,感興趣的小伙伴可以了解下
    2025-09-09
  • win10中docker部署和運行countly-server的流程

    win10中docker部署和運行countly-server的流程

    這篇文章主要記錄一下windows10中使用docker容器安裝和部署countly-server的整個流程,本文給大家講解的非常詳細,具有一定的參考借鑒價值,需要的朋友參考下吧
    2019-11-11
  • 用docker實現(xiàn)Redis主從配置的示例代碼

    用docker實現(xiàn)Redis主從配置的示例代碼

    在三臺服務(wù)器上用Docker部署Redis主從模式:Server1作為主節(jié)點,Server2和Server3配置為從節(jié)點并連接主節(jié)點,通過環(huán)境變量指定主IP,驗證復制狀態(tài)以確保高可用性
    2025-09-09
  • 簡簡單單使用Docker部署Confluence

    簡簡單單使用Docker部署Confluence

    本文使用的環(huán)境是docker17版本,重點給大家講解使用Docker部署Confluence的問題,本文給大家介紹的很好對大家的學習或工作具有一定的參考借鑒價值,需要的朋友參考下吧
    2021-06-06
  • Docker打包及部署項目完整步驟

    Docker打包及部署項目完整步驟

    這篇文章主要給大家介紹了關(guān)于Docker打包及部署項目的相關(guān)資料,Docker是一種容器化技術(shù),可以將應用程序及其依賴項打包成一個容器,方便在不同的環(huán)境中部署和運行,需要的朋友可以參考下
    2023-08-08

最新評論

彭阳县| 泰兴市| 金秀| 合水县| 沾益县| 苏州市| 边坝县| 广饶县| 长兴县| 安康市| 海宁市| 陆良县| 玛纳斯县| 呼图壁县| 青阳县| 岗巴县| 宜兰市| 怀柔区| 青冈县| 三门县| 沧州市| 开鲁县| 抚松县| 北安市| 巴彦县| 建阳市| 惠东县| 曲松县| 绥德县| 南川市| 寿光市| 宁海县| 大宁县| 金秀| 武威市| 东乡族自治县| 垫江县| 新干县| 襄樊市| 江山市| 吉木乃县|