天天看点

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();
    }
}