天天看點

RabbitMQ生産者和消費者Java Demo源碼-001

RabbitMQ生産者和消費者Java Demo源碼-001

  • ​​1. RabbitMQ 環境搭建​​
  • ​​2. 項目搭建​​
  • ​​3. 建立賬号​​
  • ​​4. 消息提供者 RabbitProducer​​
  • ​​4. 消費者代碼​​

1. RabbitMQ 環境搭建

2. 項目搭建

如果不想自己搭建可以使用我已經搭建好的:​​https://gitee.com/jack0240/rabbitmqdemo.git​​
RabbitMQ生産者和消費者Java Demo源碼-001
RabbitMQ生産者和消費者Java Demo源碼-001
RabbitMQ生産者和消費者Java Demo源碼-001
RabbitMQ生産者和消費者Java Demo源碼-001
RabbitMQ生産者和消費者Java Demo源碼-001
RabbitMQ生産者和消費者Java Demo源碼-001

設定maven, 不然下載下傳包加載很慢

​maven下載下傳配置

RabbitMQ生産者和消費者Java Demo源碼-001
添加maven配置
<!-- 日志包 -->
<dependency>
   <groupId>org.slf4j</groupId>
    <artifactId>slf4j-simple</artifactId>
    <version>1.7.30</version>
</dependency>

<dependency>
    <groupId>ch.qos.logback</groupId>
    <artifactId>logback-classic</artifactId>
    <version>1.2.3</version>
</dependency>
<!-- rabbitmq 連接配接包 -->
<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.11.0</version>
</dependency>      
RabbitMQ生産者和消費者Java Demo源碼-001
<?xml version="1.0" encoding="UTF-8"?>
<configuration debug="false">
        <!--定義日志檔案的存儲位址 勿在 LogBack 的配置中使用相對路徑-->
        <property name="LOG_HOME" value="/logs" />
            <!-- 控制台輸出 -->
        <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
            <encoder class="ch.qos.logback.classic.encoder.PatternLayoutEncoder">
            <!--格式化輸出:%d表示日期,%thread表示線程名,%-5level:級别從左顯示5個字元寬度%msg:日志消息,%n是換行符-->
            <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
            </encoder>
        </appender>
        <!-- 按照每天生成日志檔案 -->
        <appender name="FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
            <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
                <!--日志檔案輸出的檔案名-->
                <FileNamePattern>${LOG_HOME}/rabbitmqdemo.log.%d{yyyy-MM-dd}.log</FileNamePattern>
                <!--日志檔案保留天數-->
                <MaxHistory>30</MaxHistory>
            </rollingPolicy>
            <encoder class="ch.qos.logback.classic.encoder.PatternLayoutEncoder">
                <!--格式化輸出:%d表示日期,%thread表示線程名,%-5level:級别從左顯示5個字元寬度%msg:日志消息,%n是換行符-->
                <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
            </encoder>
            <!--日志檔案最大的大小-->
            <triggeringPolicy class="ch.qos.logback.core.rolling.SizeBasedTriggeringPolicy">
                <MaxFileSize>10MB</MaxFileSize>
            </triggeringPolicy>
        </appender>

        <!-- 日志輸出級别 -->
        <root level="INFO">
            <appender-ref ref="STDOUT" />
        </root>
</configuration>      

項目結構:

RabbitMQ生産者和消費者Java Demo源碼-001

3. 建立賬号

預設情況下,通路RabbitMQ服務的使用者名和密碼都是​

​guest​

​​,這個賬号有限制,預設隻能通過本地網絡(如localhost)通路,遠端網絡通路受限,是以在實作生産和消費消息之前,需要另外添加一個使用者,并設定相應的通路權限。

需要進入安裝的目錄設定:

​​

​D:\Program Files\RabbitMQ Server\rabbitmq_server-3.8.14\sbin>​

# 1.添加使用者密碼
rabbitmqctl add_user jack Jack2021

# 2.為使用者添加所有權限
rabbitmqctl set_permissions -p / jack ".*" ".*" ".*"

# 3.設定使用者為管理者
rabbitmqctl set_user_tags jack administrator      
RabbitMQ生産者和消費者Java Demo源碼-001
​​http://localhost:15672/#/users​​
RabbitMQ生産者和消費者Java Demo源碼-001

4. 消息提供者 RabbitProducer

package com.jack.rabbitmq.demo.day1;

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

/**
 * RabbitMQ 服務提供者
 * @author Jack魏
 */
public class RabbitProducer {
    private static final Logger logger = LoggerFactory.getLogger(RabbitProducer.class);

    private static final String EXCHANGE_NAME = "exchange_demo";
    private static final String ROUTING_KEY = "routingkey_demo";
    private static final String QUEUE_NAME = "queue_demo";
    private static final String IP_ADDRESS = "127.0.0.1";
    private static final int PORT = 5672;

    public static void main(String[] args) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost(IP_ADDRESS);
        factory.setPort(PORT);
        factory.setUsername("jack");
        factory.setPassword("Jack2021");
        // 建立連接配接
        Connection connection = factory.newConnection();
        // 建立信道
        Channel channel = connection.createChannel();
        // 建立一個type="direct"、持久化的、非自動删除的交換器
        channel.exchangeDeclare(EXCHANGE_NAME, "direct", true, false, null);
        // 建立一個持久化、非排他的、非自動删除的隊列
        channel.queueDeclare(QUEUE_NAME, true, false, false, null);
        // 将交換器和隊列通過路由綁定
        channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
        // 發送一條持久化的消息:hello, RabbitMQ
        String message = "hello, RabbitMQ!";
        channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY,
                MessageProperties.PERSISTENT_TEXT_PLAIN,
                message.getBytes());
        logger.info("發送資訊---{}", message);
        // 關閉資源
        channel.close();
        connection.close();
    }
}      

點選運作:

RabbitMQ生産者和消費者Java Demo源碼-001
檢視web端變化

4. 消費者代碼

package com.jack.rabbitmq.demo.day1;

import com.rabbitmq.client.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

/**
 * RabbitMQ 消費者
 * @author Jack魏
 */
public class RabbitConsumer {
    private static final Logger logger = LoggerFactory.getLogger(RabbitProducer.class);

    private static final String QUEUE_NAME = "queue_demo";
    private static final String IP_ADDRESS = "127.0.0.1";
    private static final int PORT = 5672;

    public static void main(String[] args) throws IOException, TimeoutException, InterruptedException {
        Address[] addresses = new Address[] {
                new Address(IP_ADDRESS, PORT)
        };
        ConnectionFactory factory = new ConnectionFactory();
        factory.setUsername("jack");
        factory.setPassword("Jack2021");
        // 這裡的連接配接方式與生産者的demo略有不同,注意區分
        // 建立連接配接
        Connection connection = factory.newConnection(addresses);
        // 建立信道
        final Channel channel = connection.createChannel();
        // 設定用戶端最多接受未被ack的消息的個數
        channel.basicQos(64);
        Consumer consumer = new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties properties, byte[] body) throws IOException {
                logger.info("接受到的資訊: " + new String(body));
                try {
                    TimeUnit.SECONDS.sleep(1);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                channel.basicAck(envelope.getDeliveryTag(), false);
            }
        };
        channel.basicConsume(QUEUE_NAME, consumer);
        // 等待回調函數執行完畢後,關閉資源
        TimeUnit.SECONDS.sleep(5);
        channel.close();
        connection.close();
    }
}