Day10 MQ产品选择与Rocket MQ搭建

MQ简介

Message Queue 消息队列,简称MQ,是一种异步通信机制

常见用途

  • 解耦
    • 生产者发送消息后立即返回,消费者在合适时机处理消息,减少服务之间的影响
  • 削峰
    • 高并发场景下暂存请求,平滑高峰流量,避免系统过载
  • 异步
    • 将不需要立即处理的任务放入消息队列中异步执行,减少用户请求和响应时间

消息队列的两种模式

  • 点对点模式
    • 一个生产者对应一个消费者
    • 消费者主动拉取数据,消息收到后清除消息
  • 发布订阅模式
    • 可以有多个topic
    • 消费者消费数据后不删除数据
    • 每个消费者相互独立,都可以消费到数据

常见MQ产品对比

大吞吐量大数据优先选Kafka,灵活低延迟选RabbitMQ,可靠事物国产化选RocketMQ

对比维度RabbitMQKafkaRocketMQ
开源背景Pivotal(VMware)LinkedIn → Apache阿里巴巴 → Apache
开发语言ErlangScala + JavaJava
核心模型Exchange + Queue(灵活路由)Topic + Partition(日志流)Topic + MessageQueue(轻量注册中心)
协议支持AMQP 0.9.1/1.0, MQTT, STOMP自定义二进制协议自定义协议(部分兼容 OpenMessaging)
吞吐能力1–3 万 QPS50 万+ QPS(高吞吐)10–50 万 QPS(可调优至百万级)
端到端延迟50μs – 1ms(低延迟)5–20ms(批量优化后)1–10ms
消息持久化支持(需显式开启)天然持久(顺序写磁盘)天然持久(CommitLog 顺序写)
复制机制镜像队列(主从同步)ISR 副本(In-Sync Replicas)主从同步(支持同步/异步刷盘 + 复制)
一致性保证强一致(镜像队列)最终一致(ACK=all 可近似强一致)强一致(金融级设计)
可靠性(默认)中(依赖配置)
顺序消息单队列内有序分区内有序全局有序 / 分区有序
事务消息AMQP 事务(性能差)生产者事务(仅原子性)完整二阶段事务消息
延迟消息插件支持(18 级固定延迟)不支持(需外部实现)原生支持(18 级 + 精确延迟)
死信队列支持(通过重试 Topic 模拟)支持
消息过滤Exchange 路由消费组级别Tag + SQL92 表达式
消息轨迹需插件)部分(依赖第三方)原生支持
多语言客户端极好(官方支持全)极好(社区丰富)偏 Java(其他语言生态较弱)
管理界面内置 Web UI(优秀)(依赖 Kafka Manager/UI for Kafka 等)RocketMQ Console(功能完善)
大数据生态集成极强(Flink/Spark/ELK 等)中(阿里系生态强)
动态扩容困难(集群结构固定)容易(加 Broker/Partition)较容易(加 Broker)
典型适用场景企业内部系统、任务队列、IoT、微服务通信日志采集、流处理、大数据管道、事件溯源电商交易、金融支付、订单系统、国产化项目
主要缺点吞吐低、Erlang 扩展难、集群扩容不灵活延迟高、无原生延迟/事务、运维复杂多语言支持弱、社区生态略逊于 Kafka

RocktMQ

阿里巴巴开源的使用java开发的消息中间件

单机服务搭建

bash
复制代码
# 下载官方二进制包 wget https://archive.apache.org/dist/rocketmq/5.1.4/rocketmq-all-5.1.4-bin-release.zip # 如果没有的话需要下载unzip,有的话跳过这一步 yum -y install unzip # 解压(必须用 unzip) unzip rocketmq-all-5.1.4-bin-release.zip mv rocketmq-all-5.1.4-bin-release /opt/rocketmq # 由于我的服务器内存只有2G,所以必须把内存改小 # 修改NameServer内存 vim /opt/rocketmq/bin/runserver.sh # 将 # JAVA_OPT="${JAVA_OPT} -server -Xms4g -Xmx4g -Xmn2g -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -XX:MetaspaceSize=64m -XX:MaxMetaspaceSize=128m" # 开发测试自己玩够用 # 修改Broker内存 vim /opt/rocketmq/bin/runbroker.sh # JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms512m -Xmx512m" # 修改jvm启动内存 vim /opt/rocketmq/bin/tools.sh # JAVA_OPT="${JAVA_OPT} -server -Xms1g -Xmx1g -Xmn256m -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms128m -Xmx128m -XX:MetaspaceSize=64m -XX:MaxMetaspaceSize=128m" # 修改环境变量 export NAMESRV_ADDR='localhost:9876' # 启动 NameServer(默认端口 9876) nohup sh /opt/rocketmq/bin/mqnamesrv > /var/log/rocketmq/namesrv.log 2>&1 & # 启动 Broker(连接本地 NameServer) nohup sh /opt/rocketmq/bin/mqbroker -n localhost:9876 > /var/log/rocketmq/broker.log 2>&1 & # 查看是否启动成功 jps # BrokerStartup # NamesrvStartup # 测试发送消息 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer # 测试接收消息 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer # 消息能被消费,无报错即正常

image.png

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP