ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

IM系统 —— 消息队列解耦与削峰填谷

IM系统 —— 消息队列解耦与削峰填谷 前言IM 系统中消息的转发与存储往往效率不匹配当高并发实时通信的情况下上层的大量流量很可能将下层存储压垮为解决这种现象我们引入了消息队列该项目选择 RabbitMQ将上层转发与下层存储解耦隔离以此有效提高消息管理的可靠性。一、Rabbitmq 原理要引入一个中间件我们必先了解其原理再嵌入项目中下面用一张图快速了解它的设计架构根据上图我们能理解rabbitmq 的工作流程1消费者提前声明交换机、队列并使用 Binding key 绑定交换机与队列作为后续操作的唯一标识。2消费者订阅一个队列并使用 tag 标识一个消费者订阅和一个消息投递。3生产者发布消息并携带 Routing key 匹配 Binding key 绑定相应队列此处涉及交换机的类型直接交换、广播交换、主题交换此项目选择直接交换其绑定规则要求Routing key与 Binding key完全匹配。4生产者绑定成功后将消息发布到指定队列中供消费者消费。二、二次封装根据上述原理我们能将 rabbitmq 再次封装将下层逻辑屏蔽让项目调用时更便捷1封装事件驱动模型 loopTCP连接connection客户端信道channel2声明交换机、队列并绑定交换机与队列3向指定交换机发布消息4订阅指定消息队列代码示例MQClient(const std::string user root, const std::string passwd password, const std::string host 127.0.0.1:5672) { _loop EV_DEFAULT; _handler std::make_uniqueAMQP::LibEvHandler(_loop); std::string url amqp:// user : passwd host /; AMQP::Address address(url); _conn std::make_uniqueAMQP::TcpConnection(_handler.get(), address); _channel std::make_uniqueAMQP::TcpChannel(_conn.get()); // 主线程是负责业务处理的但ev事件监控框架的启动是一个阻塞接口所以再用一个线程专门执行 _loop_thread std::thread([this](){ ev_run(_loop, 0); }); } // 声明交换机、队列并绑定 void DeclareComponents(const std::string exchange_name, const std::string queue, const std::string routine_key routine_key, AMQP::ExchangeType exchange_type AMQP::ExchangeType::direct) { // 声明交换机 _channel-declareExchange(exchange_name, exchange_type) .onError([exchange_name](const char *message){ LOG_ERROR(声明交换机 {} 失败 {}, exchange_name, message); exit(0); }) .onSuccess([exchange_name](){ //返回值Deferred LOG_INFO(声明交换机 {} 成功, exchange_name); }); // 声明队列 _channel-declareQueue(queue) .onError([queue](const char *message){ LOG_ERROR(声明队列 {} 失败 {}, queue, message); exit(0); }) .onSuccess([queue](){ LOG_INFO(声明队列 {} 成功, queue); }); // 交换机绑定队列 _channel-bindQueue(exchange_name, queue, routine_key) .onError([exchange_name, queue](const char *message){ LOG_ERROR(交换机{}-队列{}绑定失败 {}, exchange_name, queue, message); exit(0); }) .onSuccess([exchange_name, queue](){ LOG_INFO(交换机{}-队列{}绑定成功, exchange_name, queue); }); } // 发布消息至交换机 bool Publish(const std::string exchange_name, const std::string msg, const std::string routine_key routine_key) { LOG_DEBUG(向交换机 {}-{} 发布消息, exchange_name, routine_key); if(!_channel-publish(exchange_name, routine_key, msg)) { LOG_ERROR({} 发布失败, msg); return false; } return true; } // 订阅队列 bool Consume(const std::string queue, const MessageCallback cb, const std::string tag consume_tag) { LOG_DEBUG(开始订阅 {} 队列消息, queue); _channel-consume(queue, tag) .onReceived([this, cb](const AMQP::Message message, uint64_t deliveryTag, bool redelivered){ cb(message.body(), message.bodySize()); // 消费 _channel-ack(deliveryTag); // 确认完成消费 }) .onError([queue](const char *message){ LOG_ERROR(订阅队列 {} 失败 {}, queue, message); return false; }); return true; }注意由于消息队列的 ev_run 事件监控接口为阻塞接口我们需要使用独立线程监听所以在MQClient 的析构函数中需要主线程发送异步通知让线程停止事件监控。ev_async_init(ev_async *void (*)(struct ev_loop *, ev_async *, int));ev_async_start(struct ev_loop *ev_async *);ev_async_send(struct ev_loop * ev_async *);接口解析1绑定异步观察器与回调函数。2将异步观察器注册到事件监控的监控列表中。3主线程唤醒阻塞在epoll_wait上的事件循环循环醒来后执行回调。tips异步通知之后避免使用类似 ev_loop_destroy 的接口异步通知只是通知从线程执行销毁操作如果在从线程执行该操作之前使用 ev_loop_destroy 会导致从线程提前销毁从而引发二次析构。
返回列表