欢迎光临
我们一直在努力

【框架工具#9】RabbitMQ 安装和使用

在这里插入图片描述

📃个人主页:island1314

⛺️ 欢迎关注:👍点赞 👂🏽留言 😍收藏 💞 💞 💞

  • 生活总是不会一帆风顺,前进的道路也不会永远一马平川,如何面对挫折影响人生走向 – 《人民日报》

🔥 目录

  • RabbitMQ 安装和使用
    • 一、基本概述
      • 1. 核心概念(AMQP 模型)
      • 2. Exchange 交换机类型
      • 3. 消息可靠性机制
      • 4. 高级特性
      • 5. 工作模式(快速记忆)
      • 6. 性能与扩展
      • 7. 与 Kafka 的简单对比
      • 8. 典型使用场景
    • 二、理解 RabbitMQ 相关核心概念
    • 三、安装
      • 1. 安装教程
      • 2. 简单使用
      • 3. 安装 RabbitMQ 的 C++客户端库
    • 三、常用类和接口
      • 1. Channel类
      • 2. libev
    • 四、基本使用
      • 1. consumer
      • 2. publish
    • 五、二次封装
      • 1. 实现思路
      • 2. 具体实现
      • 3. 封装测试

RabbitMQ 安装和使用

一、基本概述

RabbitMQ 是一个开源的 消息代理(message broker),用 Erlang 编写,实现了 AMQP(Advanced Message Queuing Protocol) 标准,并支持 MQTT、STOMP 等多种协议。它本质上是一个“邮局”:生产者把消息发到 RabbitMQ,RabbitMQ 把消息投递给消费者,二者无需同时在线,实现应用解耦、削峰填谷、异步处理。

1. 核心概念(AMQP 模型)

概念说明
Producer 发消息的应用。
Consumer 收消息的应用。
Queue 真正存储消息的缓冲区,FIFO(可配置优先级、TTL 等)。
Exchange 接收生产者消息,根据 路由规则 把消息投放到 0~N 个队列。
Binding 队列与交换器之间的绑定关系,带一个 routing key。
Routing Key 生产者发送消息时带的“地址”,Exchange 用它 + Binding 决定消息去哪。
Virtual Host (vhost) 逻辑隔离单元,类似数据库里的“库”,可配独立权限、队列、交换器。

2. Exchange 交换机类型

enum ExchangeType
{
fanout,
direct,
topic,
headers,
consistent_hash,
message_deduplication
};

四种常见交换机如下:

类型路由规则典型场景
direct 精确匹配 routing key。 单播任务分发。
topic 模糊匹配(点分单词,* 单层 # 多层)。 多类目日志订阅。
fanout 无视 key,广播到所有绑定队列。 群发通知、配置刷新。
headers 根据消息头属性匹配,性能差,很少用。 复杂头匹配场景。

3. 消息可靠性机制

  • 生产者确认(Publisher Confirm):异步 ACK/NACK,告知消息是否到达 Exchange/Queue。
  • 事务机制(txSelect/txCommit)已废弃,性能差。
  • 持久化:队列、消息、交换器都可标记为 durable,重启后自动恢复。
  • 消费者手动 ack:默认自动 ack 会丢失消息;改为手动 basic.ack 只有处理成功才确认。
  • Publisher Returns:不可路由消息会回发 basic.return,可做兜底告警。
  • 4. 高级特性

    • TTL(消息/队列级别生存时间)
    • 死信队列 DLX(TTL 过期、reject、队列满后转发到 DLX)
    • 延迟队列(TTL + DLX 或插件 rabbitmq_delayed_message_exchange)
    • 优先级队列(x-max-priority)
    • 队列镜像(高可用,3.8 后推荐 Quorum Queue 替代传统镜像)
    • 流式队列(RabbitMQ 3.9+,类似 Kafka 的 log,百万级吞吐)
    • Federation/Shovel插件:跨机房、跨集群消息复制
    • 管理界面 (rabbitmq_management):Web UI、REST API、实时监控
    • Prometheus+Grafana 导出器:指标全覆盖

    5. 工作模式(快速记忆)

  • 简单队列(Hello World):默认 Exchange,routing key = 队列名。
  • Work Queues:多个消费者竞争消费,round-robin 分发,可配 prefetch=1 实现“能者多劳”。
  • Publish/Subscribe(fanout)
  • Routing(direct)
  • Topics(topic)
  • RPC(reply_to + correlation_id 实现异步 RPC)
  • 6. 性能与扩展

    • 单机 万级~十万级 TPS(消息大小 1 KB、持久化、confirm 模式下)。
    • 队列绑定 CPU 线程模型(Erlang 轻量进程),横向扩容加节点即可。
    • 3.8 引入 Quorum Queue(Raft 协议)替代镜像队列,更高吞吐、更低延迟、更强一致性。
    • 流式队列(RabbitMQ Streams)可替代 Kafka 部分场景,百万级 QPS。

    7. 与 Kafka 的简单对比

    RabbitMQKafka
    定位 通用消息代理,路由丰富 高吞吐分布式日志
    模型 队列 + Exchange 分区日志
    消息保留 consumer ack 后删除 按时间/大小保留,可重放
    吞吐量 单机 万~十万 单机 百万+
    延迟 毫秒级 毫秒~秒级
    路由 灵活(topic、header、fanout) 简单 key/hash
    事务/顺序 单队列顺序 分区内顺序

    8. 典型使用场景

    • 订单系统:下单后异步扣库存、发优惠券、发短信。
    • 秒杀削峰:把瞬间请求转成消息,平滑消费。
    • 微服务解耦:AB 服务不直接调用,通过消息协作。
    • 日志/审计:多系统把日志发到统一队列,集中处理。
    • 延迟任务:订单 30 分钟未支付自动关闭(TTL+DLX)。

    一句话总结:RabbitMQ = 轻量级、高可靠、路由灵活的消息中间件,在“需要复杂路由、消息不丢、延迟低”的场景下,是 Kafka 之外最稳妥的选择。

    二、理解 RabbitMQ 相关核心概念

    Exchange(交换机)、Queue(队列)、Binding(绑定)、Message(消息) 就是 RabbitMQ 的四大金刚。

    image-20250916204417405

    把 RabbitMQ 想成“快递公司”,四样东西立刻有画面:

  • Message = 快递包裹 业务数据 + 标签(routing key、headers 等)
  • Queue = 小区快递柜 真正放包裹的地方,消费者从柜子里取。
  • Exchange = 分拨中心 快递员(生产者)只能把包裹扔进分拨中心,不能直接塞进柜子。
  • Binding = 分拨规则 “分拨中心 → 哪些柜子”的路由表,决定包裹最后落到哪个柜子。
  • 生活例子:外卖订单

    我们点了一份奶茶和一份汉堡,平台(生产者)把订单拆成两条消息:

    • 消息 A:{"item":"奶茶"},routing key = drink
    • 消息 B:{"item":"汉堡"},routing key = food

    Exchange 是“商家出餐调度中心”,类型 = direct(精确匹配)。

    店里有两个Queue(快递柜):

    队列名绑定键(Binding key)后台消费者
    q_drink drink 奶茶师
    q_food food 炸锅师

    流程走一遍:

  • 调度中心收到消息 A,看 routing key = drink,查路由表 → 命中 q_drink → 消息 A 进柜。
  • 消息 B 同理进 q_food。
  • 奶茶师、炸锅师各自从自己柜子取单消费。
  • 如果又来一条 {"item":"薯条"},routing key = food,一样会落入 q_food,不需要新建队列; 若 routing key = snack,没有对应绑定,包裹就被丢弃(或进死信)。

    一句话总结

    • 生产者 → Exchange → Binding → Queue → 消费者
    • 你把“路由键”填好,剩下的“分拨-落地-取件”全由 RabbitMQ 完成。

    三、安装

    1. 安装教程

    sudo apt install rabbitmq-server

    2. 简单使用

    # 启动服务
    sudo systemctl start rabbitmq-server.service
    # 查看服务状态
    sudo systemctl status rabbitmq-server.service

    注意:安装完成的时候默认有个用户 guest ,但是权限不够,要创建一个 administrator 用户,才可以做为远程登录和发表订阅消息:

    # 添加用户
    lighthouse@VM-8-10-ubuntu:~$ sudo rabbitmqctl add_user root 123456
    Adding user "root" ...
    Done. Don't forget to grant the user permissions to some virtual hosts! See 'rabbitmqctl help set_permissions' to learn more.
    # 上面提示成功创建了一个 root 用户,但此时它还什么都不能做——默认没有任何虚拟主机(vhost)的访问权限。下一步就是给它授权

    # 设置用户 tag
    lighthouse@VM-8-10-ubuntu:~$ sudo rabbitmqctl set_user_tags root administrator
    Setting tags for user "root" to [administrator] ...
    # 设置用户权限
    lighthouse@VM-8-10-ubuntu:~$ sudo rabbitmqctl set_permissions -p / root "." "." ".*"
    Setting permissions for user "root" in vhost "/" ...

    # RabbitMQ 自带了 Web 管理界面 通过下面指令开启
    lighthouse@VM-8-10-ubuntu:~$ sudo rabbitmq-plugins enable rabbitmq_management
    Enabling plugins on node rabbit@VM-8-10-ubuntu:
    rabbitmq_management
    The following plugins have been configured:
    rabbitmq_management
    rabbitmq_management_agent
    rabbitmq_web_dispatch
    Applying plugin configuration to rabbit@VM-8-10-ubuntu...
    The following plugins have been enabled:
    rabbitmq_management
    rabbitmq_management_agent
    rabbitmq_web_dispatch

    started 3 plugins.

    而且其默认端口是 15672,如下:

    lighthouse@VM-8-10-ubuntu:~$ netstat -anptu | grep 15672
    (Not all processes could be identified, non-owned process info
    will not be shown, you would have to be root to see it all.)
    tcp 0 0 0.0.0.0:15672 0.0.0.0:* LISTEN –

    访问 webUI 界面,结果如下:(网站为:ip:15672)

    image-20250916212001182

    输入账号密码,显示如下:

    image-20251228171906023

    3. 安装 RabbitMQ 的 C++客户端库

    • C 语言库:https://github.com/alanxz/rabbitmq-c
    • C++ 库: https://github.com/CopernicaMarketingSoftware/AMQP-CPP/tree/master

    我们这里使用 AMQP-CPP 库来编写客户端程序

    安装 AMQP-CPP

    sudo apt install libev-dev #libev 网络库组件
    git clone https://github.com/CopernicaMarketingSoftware/AMQP-CPP.git

    cd AMQP-CPP/

    make & make install

    安装报错解决(Make 的时候可能会出现如下报错)

    /usr/include/openssl/macros.h:147:4: error: #error "OPENSSL_API_COMPAT expresses an impossible API compatibility level"
    147 | # error "OPENSSL_API_COMPAT expresses an impossible API compatibility level"
    | ^~~~~
    In file included from /usr/include/openssl/ssl.h:18,
    from linux_tcp/openssl.h:20,
    from linux_tcp/openssl.cpp:12:
    /usr/include/openssl/bio.h:687:1: error: expected constructor, destructor, or type conversion before ‘DEPRECATEDIN_1_1_0’

    这种错误,表示 ssl 版本出现问题。 解决方案:卸载当前的 ssl 库,重新进行修复安装

    lighthouse@VM-8-10-ubuntu:AMQP-CPP$ dpkg -l |grep ssl
    ii erlang-ssl 1:24.2.1+dfsg-1ubuntu0.5
    ii libevent-openssl-2.1-7:amd64 2.1.12-stable-1build3
    ii libgnutls-openssl27:amd64 3.7.3-4ubuntu1.7
    ii libssl-dev:amd64 3.0.2-0ubuntu1.19
    ii libssl3:amd64 3.0.2-0ubuntu1.19
    ii libxmlsec1-openssl:amd64 1.2.33-1build2
    ii libzstd-dev:amd64 1.4.8+dfsg-3build1
    ii libzstd1:amd64 1.4.8+dfsg-3build1
    ii openssl 3.0.2-0ubuntu1.18
    ii python3-brotli 1.0.9-2build6
    ii python3-openssl 21.0.0-1
    ii zstd 1.4.8+dfsg-3build1

    修复安装指令如下:

    sudo dpkg -P –force-all libevent-openssl-2.1-7
    sudo dpkg -P –force-all openssl
    sudo dpkg -P –force-all libssl-dev

    sudo apt –fix-broken install

    修复之后 重新 make 即可

    三、常用类和接口

    1. Channel类

    Channel通过连接建立,用于执行所有的队列操作。每个Channel对象代表了与RabbitMQ服务器的一个 虚拟连接,其可以独立进行消息的发送、接收、确认等操作。

    常用接口

    • AMQP::Channel:连接的工作单位,通过它你可以定义交换机(exchange)、队列(queue)和绑定(binding)等;通过这个类可以实现发送和接收消息
    • AMQP::Connection:表示与 RabbitMQ 服务器的连接。每个连接对应一个与 RabbitMQ 服务器的 TCP 连接
    • AMQP::Queue:表示队列。你可以将消息发布到队列中,或从队列中消费消息
    • AMQP::Exchange:用于消息路由。通过交换机,你可以决定将消息发送到哪个队列
    • AMQP::Message:消息类,表示通过 RabbitMQ 发送的消息。它包括消息体、属性和标头等

    2. libev

    其是一个高性能事件循环库,通常用于处理 I/O 多路复用,支持定时器、信号和 I/O 事件,项目中主要用于构建网络应用程序,处理大量并发连接

    常用接口

    • ev_loop:事件循环的核心对象。在 libev 中,所有的事件和回调都注册到 ev_loop 中,程序的执行依赖于它不断地检查和处理这些事件
    • ev_io:用于监控 I/O 事件(如可读/可写)。通常与套接字相关,用于处理来自网络的事件
    • ev_timer:用于定时事件,类似于定时器,用于在指定的时间间隔后触发回调
    • ev_signal:用于捕获信号,在接收到信号时,触发回调函数

    四、基本使用

    Makefile 文件 如下:

    all: publish consume
    publish: publish.cc
    g++ -o $@ $^ -lamqpcpp -lev -lfmt -lspdlog -lgflags -g
    consume: consume.cc
    g++ -o $@ $^ -lamqpcpp -lev -lfmt -lspdlog -lgflags -g
    .PHONY:clean
    clean:
    rm -f publish consume

    注意:RabbitMQ 客户端的端口不是 15672,15672 是其 UI 端口

    lighthouse@VM-8-10-ubuntu:rabbitmq$ cat /etc/rabbitmq/rabbitmq-env.conf
    # Defaults to rabbit. This can be useful if you want to run more than one node
    # per machine – RABBITMQ_NODENAME should be unique per erlang-node-and-machine
    # combination. See the clustering on a single machine guide for details:
    # http://www.rabbitmq.com/clustering.html#single-machine
    #NODENAME=rabbit

    # By default RabbitMQ will bind to all interfaces, on IPv4 and IPv6 if
    # available. Set this if you only want to bind to one network interface or#
    # address family.
    #NODE_IP_ADDRESS=127.0.0.1

    # Defaults to 5672.
    #NODE_PORT=5672

    1. consumer

    基本逻辑

    1. 初始化事件循环:通过 libev 创建一个事件循环 loop,并使用 AMQP::LibEvHandler 将事件循环与 AMQP 客户端关联

    • 快速理解:loop类似于一个活动的组织者负责每个成员的IO,handler则是助手(任务:将所有信息和任务连接在一起)

    2. 建立与 RabbitMQ 的连接:通过 AMQP::TcpConnection 建立与 RabbitMQ 服务器的 TCP 连接

    • 理解:类似于活动的接待者(Connection)为来者提供了直接交流通道,这个通道(Address)可以和别人发送消息

    3. 声明交换机和队列:声明交换机(test-exchange)和队列(test-queue),并绑定它们

    • 理解:声明你的谈话方式和交流位置

    4. 消费消息:订阅队列中的消息并注册回调函数 MessageCb 来处理消息

    • 理解:当被人给你传递消息后你是如何做出反应

    5. 启动事件循环:调用 ev_run 启动 libev 事件循环,等待事件并调用回调函数处理异步消息

    代码实现

    #include <ev.h>
    #include <amqpcpp.h>
    #include <amqpcpp/libev.h>
    #include <openssl/ssl.h>
    #include <openssl/opensslv.h>
    #include <iostream>
    #include <string>

    // 消息回调处理函数的实现
    void MessageCb(AMQP::TcpChannel* channel, const AMQP::Message &message, uint64_t deliveryTag, bool redelivered){
    std::string msg;
    msg.assign(message.body(), message.bodySize());
    std::cout << msg << std::endl;
    channel->ack(deliveryTag); // 对消息进行确认
    }

    int main()
    {
    // 1. 实例化底层网络通信框架的I/O事件监控句柄
    auto *loop = EV_DEFAULT;
    // 2. 实例化libEvHandler句柄 — 将AMQP框架与事件监控关联起来
    AMQP::LibEvHandler handler(loop);
    // 2.5. 实例化连接对象
    AMQP::Address host("amqp://root:123456@127.0.0.1:5672/");
    AMQP::TcpConnection connection(&handler, host); // TCP 连接
    // 3. 实例化信道对象
    AMQP::TcpChannel channel(&connection);
    // 4. 声明交换机
    channel.declareExchange("test-exchange", AMQP::ExchangeType::direct)
    .onError([](const char* msg){
    std::cout << "声明交换机失败: " << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange 交换机创建成功\\n";
    });
    // 5. 声明队列
    channel.declareQueue("test-queue")
    .onError([](const char* msg){
    std::cout << "声明队列失败: " << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange 队列创建成功\\n";
    });
    // 6. 针对交换机和队列进行绑定 并指定路由键(Routing Key)
    channel.bindQueue("test-exchange", "test-queue", "test-queue-key")
    .onError([](const char* msg){
    std::cout << "test-exchange – test-queue 绑定失败:" << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange – test-queue 绑定成功\\n";
    });
    // 7. 订阅队列消息 — 设置消息回调处理函数
    auto callback = std::bind(MessageCb, &channel, std::placeholders::_1, std::placeholders::_2, std::placeholders::_3);
    channel.consume("test-queue", "consume-tag")
    .onReceived(callback)
    .onError([](const char* msg){
    std::cout << "订阅 test-queue 队列消息失败: " << msg << "\\n";
    exit(0);
    });
    // 启动底层网络通信框架–开启I/O
    ev_run(loop, 0);
    return 0;
    }

    2. publish

    逻辑梳理

  • 连接 RabbitMQ:通过 AMQP 协议连接到 RabbitMQ 服务器,设置通信所需的交换机和队列
  • 声明交换机和队列:声明一个交换机和队列,并将它们通过路由键绑定在一起
  • 发布消息:生产者向交换机发布消息,消息根据绑定的路由规则被转发到队列中
  • 事件循环:启动 libev 事件循环,处理网络 I/O 操作,确保消息的传输和处理
  • 代码实现(和 consume 大差不差,只是前者是订阅,publish 是发布)

    #include <ev.h>
    #include <amqpcpp.h>
    #include <amqpcpp/libev.h>
    #include <openssl/ssl.h>
    #include <openssl/opensslv.h>

    int main()
    {
    // 1. 实例化底层网络通信框架的I/O事件监控句柄
    auto *loop = EV_DEFAULT;
    // 2. 实例化libEvHandler句柄 — 将AMQP框架与事件监控关联起来
    AMQP::LibEvHandler handler(loop);
    // 2.5. 实例化连接对象
    AMQP::Address host("amqp://root:123456@127.0.0.1:5672/");
    AMQP::TcpConnection connection(&handler, host); // TCP 连接
    // 3. 实例化信道对象
    AMQP::TcpChannel channel(&connection);
    // 4. 声明交换机
    channel.declareExchange("test-exchange", AMQP::ExchangeType::direct)
    .onError([](const char* msg){
    std::cout << "声明交换机失败: " << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange 交换机创建成功\\n";
    });
    // 5. 声明队列
    channel.declareQueue("test-queue")
    .onError([](const char* msg){
    std::cout << "声明队列失败: " << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange 队列创建成功\\n";
    });
    // 6. 针对交换机和队列进行绑定 并指定路由键(Routing Key)
    channel.bindQueue("test-exchange", "test-queue", "test-queue-key")
    .onError([](const char* msg){
    std::cout << "test-exchange – test-queue 绑定失败:" << msg << std::endl;
    exit(0);
    })
    .onSuccess([](){
    std::cout << "test-exchange – test-queue 绑定成功\\n";
    });
    // 7. 向交换机发布消息
    for(int i = 0; i < 10; ++i){
    std::string msg = "Hello Bite-" + std::to_string(i);
    bool ret = channel.publish("test-exchange", "test-queue-key", msg); // 发布消息
    if (!ret) {
    std::cout << "publish 失败!\\n"; // 发布消息失败的回调
    }
    }
    // 8. 启动底层网络通信框架–开启I/O
    ev_run(loop, 0);
    return 0;
    }

    运行结果如下:

    lighthouse@VM-8-10-ubuntu:rabbitmq$ ./publish
    test-exchange 交换机创建成功
    test-exchange 队列创建成功
    test-exchange – test-queue 绑定成功

    lighthouse@VM-8-10-ubuntu:rabbitmq$ ./consume
    test-exchange 交换机创建成功
    test-exchange 队列创建成功
    test-exchange – test-queue 绑定成功
    Hello Bite-0
    Hello Bite-1
    Hello Bite-2
    Hello Bite-3
    Hello Bite-4
    Hello Bite-5
    Hello Bite-6
    Hello Bite-7
    Hello Bite-8
    Hello Bite-9

    五、二次封装

    1. 实现思路

    项目需求:交换机与队列直接交换,实现一台主机将消息发布给另一个主机进行处理

    封装接口

    • 提供 声明指定交换机与队列,然后绑定其功能
    • 提供向 指定交换机发布消息 功能
    • 提供 订阅指定队列消息,同时设置回调函数进行消息消费处理的功能

    2. 具体实现

    #include <ev.h>
    #include <amqpcpp.h>
    #include <amqpcpp/libev.h>
    #include <thread>
    #include <openssl/ssl.h>
    #include <openssl/opensslv.h>
    #include "logger.hpp"

    class MQClient{
    public:
    using MessageCallback = std::function<void(const char*, size_t)>;
    using ptr = std::shared_ptr<MQClient>;

    MQClient(const std::string& usr, const std::string passwd, const std::string host){
    // 1. 实例化底层网络通信框架的I/O事件监控句柄
    _loop = EV_DEFAULT;
    // 2. 实例化libEvHandler句柄
    _handler = std::make_unique<AMQP::LibEvHandler>(_loop);
    // 3. 实例化连接对象
    std::string url = "amqp://" + usr + ":" + passwd + "@" + host + "/";
    AMQP::Address address(url);
    _connection = std::make_unique<AMQP::TcpConnection>(_handler.get(), address);
    // _connection = std::make_unique<AMQP::TcpConnection>(&_handler, address);
    // 4. 实例化信道对象
    _channel = std::make_unique<AMQP::TcpChannel>(_connection.get());
    // 5. 启动底层网络通信框架–开启I/O
    _loop_thread = std::thread([this](){
    ev_run(_loop, 0);
    });
    }
    ~MQClient(){
    ev_async_init(&_async_watcher, watcher_callback);
    ev_async_start(_loop, &_async_watcher);
    ev_async_send(_loop, &_async_watcher);
    _loop_thread.join();
    _loop = nullptr;
    }

    void declareComponents(const std::string &exchange, const std::string &queue,
    const std::string &routing_key = "routing_key",
    AMQP::ExchangeType exchange_type = AMQP::ExchangeType::direct)
    {
    // 1. 声明交换机
    _channel->declareExchange(exchange, exchange_type)
    .onError([](const char* msg){
    LOG_ERROR("声明交换机失败: {}", msg);
    exit(0);
    })
    .onSuccess([exchange](){
    LOG_DEBUG("{} 交换机创建成功", exchange);
    });
    // 2. 声明队列
    _channel->declareQueue("test-queue")
    .onError([](const char* msg){
    LOG_ERROR("声明队列失败: {}", msg);
    exit(0);
    })
    .onSuccess([queue](){
    LOG_DEBUG("{} 队列创建成功", queue);
    });
    // 3. 针对交换机和队列进行绑定
    _channel->bindQueue(exchange, queue, routing_key)
    .onError([exchange, queue](const char *message) {
    LOG_ERROR("{} – {} 绑定失败:", exchange, queue);
    exit(0);
    })
    .onSuccess([exchange, queue, routing_key](){
    LOG_ERROR("{} – {} – {} 绑定成功!", exchange, queue, routing_key);
    });
    }

    bool publish(const std::string &exchange, const std::string &msg,
    const std::string &routing_key = "routing_key")
    {
    if(!_channel->publish(exchange, routing_key, msg)){
    LOG_ERROR("{} 发布消息失败:", exchange);
    return false;
    }
    return true;
    }

    void consume(const std::string& queue, const MessageCallback &cb)
    {
    LOG_DEBUG("开始订阅 {} 队列消息!", queue);
    _channel->consume(queue, "consume->tag")
    .onReceived([this, cb](const AMQP::Message &message, u_int64_t deliveryTag, bool redelivered){
    cb(message.body(), message.bodySize());
    _channel->ack(deliveryTag);
    })
    .onError([queue](const char *message){
    LOG_ERROR("订阅 {} 队列消息失败: {}", queue, message);
    exit(0);
    });
    }

    private:
    static void watcher_callback(struct ev_loop *loop, ev_async *watcher, int32_t revents) {
    ev_break(loop, EVBREAK_ALL);
    }

    struct ev_loop *_loop;
    struct ev_async _async_watcher;
    std::unique_ptr<AMQP::LibEvHandler> _handler;
    std::unique_ptr<AMQP::TcpConnection> _connection;
    std::unique_ptr<AMQP::TcpChannel> _channel;
    std::thread _loop_thread;
    };

    3. 封装测试

    consume.cc

    #include "../../../common/rabbitmq.hpp"
    #include <gflags/gflags.h>

    DEFINE_string(user, "root", "rabbitmq访问用户名");
    DEFINE_string(pswd, "123456", "rabbitmq访问密码");
    DEFINE_string(host, "127.0.0.1:5672", "rabbitmq服务器地址信息 host:port");

    DEFINE_bool(run_mode, false, "程序的运行模式,false-调试; true-发布;");
    DEFINE_string(log_file, "", "发布模式下,用于指定日志的输出文件");
    DEFINE_int32(log_level, 0, "发布模式下,用于指定日志输出等级");

    void callback(const char* body, size_t sz){
    std::string msg;
    msg.assign(body, sz);
    std::cout << msg << std::endl;
    }

    int main(int argc, char* argv[]){
    google::ParseCommandLineFlags(&argc, &argv, true);
    init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);

    MQClient client(FLAGS_user, FLAGS_pswd, FLAGS_host);
    client.declareComponents("test-exchange", "test-queue");
    client.consume("test-queue", callback);

    std::this_thread::sleep_for(std::chrono::seconds(60));
    return 0;
    }

    publish.cc

    #include "../../../common/rabbitmq.hpp"
    #include <gflags/gflags.h>

    DEFINE_string(user, "root", "rabbitmq访问用户名");
    DEFINE_string(pswd, "123456", "rabbitmq访问密码");
    DEFINE_string(host, "127.0.0.1:5672", "rabbitmq服务器地址信息 host:port");

    DEFINE_bool(run_mode, false, "程序的运行模式,false-调试; true-发布;");
    DEFINE_string(log_file, "", "发布模式下,用于指定日志的输出文件");
    DEFINE_int32(log_level, 0, "发布模式下,用于指定日志输出等级");

    int main(int argc, char *argv[])
    {
    google::ParseCommandLineFlags(&argc, &argv, true);
    init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);

    MQClient client(FLAGS_user, FLAGS_pswd, FLAGS_host);

    client.declareComponents("test-exchange", "test-queue");

    for (int i = 0; i < 10; i++) {
    std::string msg = "Hello Bite-" + std::to_string(i);
    if (!client.publish("test-exchange", msg)) {
    std::cout << "publish 失败!\\n";
    }
    }
    std::this_thread::sleep_for(std::chrono::seconds(3));
    return 0;
    }

    结果如下:

    lighthouse@VM-8-10-ubuntu:rabbitmq$ ./publish
    [default-logger][21:53:12][1573326][debug ][../../../common/rabbitmq.hpp:50] test-exchange 交换机创建成功
    [default-logger][21:53:12][1573326][debug ][../../../common/rabbitmq.hpp:59] test-queue 队列创建成功
    [default-logger][21:53:12][1573326][error ][../../../common/rabbitmq.hpp:68] test-exchange – test-queue – routing_key 绑定成功!

    lighthouse@VM-8-10-ubuntu:rabbitmq$ ./consume
    [default-logger][21:53:14][1573347][debug ][../../../common/rabbitmq.hpp:85] 开始订阅 test-queue 队列消息!
    [default-logger][21:53:14][1573349][debug ][../../../common/rabbitmq.hpp:50] test-exchange 交换机创建成功
    [default-logger][21:53:14][1573349][debug ][../../../common/rabbitmq.hpp:59] test-queue 队列创建成功
    [default-logger][21:53:14][1573349][error ][../../../common/rabbitmq.hpp:68] test-exchange – test-queue – routing_key 绑定成功!
    Hello Bite-0
    Hello Bite-1
    Hello Bite-2
    Hello Bite-3
    Hello Bite-4
    Hello Bite-5
    Hello Bite-6
    Hello Bite-7
    Hello Bite-8
    Hello Bite-9

    【★,°:.☆( ̄▽ ̄)/$:.°★ 】那么本篇到此就结束啦,如果有不懂 和 发现问题的小伙伴可以在评论区说出来哦,同时我还会继续更新相关的内容,请持续关注我 !!

    在这里插入图片描述

    赞(0)
    未经允许不得转载:171主机测评 » 【框架工具#9】RabbitMQ 安装和使用
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址