Day10 MQ产品选择与Rocket MQ搭建
MQ简介
Message Queue 消息队列,简称MQ,是一种异步通信机制
常见用途
- 解耦
- 生产者发送消息后立即返回,消费者在合适时机处理消息,减少服务之间的影响
- 削峰
- 高并发场景下暂存请求,平滑高峰流量,避免系统过载
- 异步
- 将不需要立即处理的任务放入消息队列中异步执行,减少用户请求和响应时间
消息队列的两种模式
- 点对点模式
- 一个生产者对应一个消费者
- 消费者主动拉取数据,消息收到后清除消息
- 发布订阅模式
- 可以有多个topic
- 消费者消费数据后不删除数据
- 每个消费者相互独立,都可以消费到数据
常见MQ产品对比
大吞吐量大数据优先选Kafka,灵活低延迟选RabbitMQ,可靠事物国产化选RocketMQ
| 对比维度 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 开源背景 | Pivotal(VMware) | LinkedIn → Apache | 阿里巴巴 → Apache |
| 开发语言 | Erlang | Scala + Java | Java |
| 核心模型 | Exchange + Queue(灵活路由) | Topic + Partition(日志流) | Topic + MessageQueue(轻量注册中心) |
| 协议支持 | AMQP 0.9.1/1.0, MQTT, STOMP | 自定义二进制协议 | 自定义协议(部分兼容 OpenMessaging) |
| 吞吐能力 | 1–3 万 QPS | 50 万+ 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 # 消息能被消费,无报错即正常

评论
问答助学
相关内容
0个评论
全部评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
