第 4 章 · 消息与服务治理:RabbitMQ / Kafka / RocketMQ / Nacos
4.1 RabbitMQ 3(Management)
适用场景
| 选用 | 典型业务 |
|---|---|
| 任务异步化 | 下单后发邮件、发短信 |
| 削峰填谷 | 秒杀订单写入队列,Worker 匀速消费 |
| 路由灵活 | 按 routing key 分发到不同队列 |
| 需要 可靠投递、ACK、死信队列 | 支付回调、对账任务 |
不选用:超大数据流日志管道(用 Kafka);极简 Pub/Sub 且无持久化(可用 Redis)。
商店部署参数
| 项 | 值 |
|---|---|
| id | rabbitmq |
| 镜像 | rabbitmq:3-management-alpine |
| 目录默认端口 | 15672(管理 UI,非 AMQP) |
| AMQP 端口 | 5672(需 应用部署 → 端口 补映射) |
| PVC | 10 Gi → /var/lib/rabbitmq |
环境变量
env:
- name: RABBITMQ_DEFAULT_USER
value: "app"
- name: RABBITMQ_DEFAULT_PASS
value: "强密码"
- name: RABBITMQ_DEFAULT_VHOST
value: "/"
连接语法
AMQP URI
amqp://app:密码@rabbitmq:5672/%2F
Spring AMQP
spring:
rabbitmq:
host: rabbitmq
port: 5672
username: app
password: ${RABBITMQ_PASSWORD}
virtual-host: /
Python pika
params = pika.URLParameters('amqp://app:密码@rabbitmq:5672/')
核心概念与配置语法
| 概念 | 说明 |
|---|---|
| Exchange | 接收消息并路由 |
| Queue | 存储消息 |
| Binding | Exchange 与 Queue 绑定规则 |
| Routing Key | 路由键 |
Exchange 类型
| 类型 | 路由规则 |
|---|---|
direct | routing key 完全匹配 |
fanout | 广播到所有绑定队列 |
topic | order.* 模式匹配 |
headers | 按消息头匹配 |
管理 API(HTTP,端口 15672)
# 创建 vhost
curl -u app:密码 -X PUT http://rabbitmq:15672/api/vhosts/shop
# 查看队列
curl -u app:密码 http://rabbitmq:15672/api/queues
Java 声明队列(示例)
@Bean Queue orderQueue() { return new Queue("order.created", true); }
@Bean DirectExchange orderExchange() { return new DirectExchange("order"); }
@Bean Binding binding() {
return BindingBuilder.bind(orderQueue()).to(orderExchange()).with("order.created");
}
部署后必做(小紫)
- 应用部署 → 端口 → 添加
5672NodePort - 管理 UI:
http://<节点IP>:<15672的NodePort> - 节点防火墙 放行
小紫运维
队列堆积 → 管理 UI Queues 看 Ready 数;消费慢 → 扩容 Worker 副本(非 RabbitMQ 本身);磁盘告警 → 存储 扩容。
4.2 Kafka 3.7(Bitnami)
适用场景
| 选用 | 场景 |
|---|---|
| 日志采集、埋点管道 | 高吞吐写入 |
| 事件溯源、CDC | 多订阅者重复读 |
| 流式处理(Flink/Spark 接入) | 分区并行 |
不选用:低吞吐、路由复杂的小任务(RabbitMQ 更简单);团队无运维 Kafka 经验时的首选用 MQ。
商店部署参数
| 项 | 值 |
|---|---|
| id | kafka |
| 镜像 | bitnami/kafka:3.7 |
| 端口 | 9092 |
| PVC | 20 Gi → /bitnami/kafka |
| 资源 | 中型 · 2 核 4Gi |
环境变量(单节点 KRaft 模式)
env:
- name: KAFKA_CFG_NODE_ID
value: "0"
- name: KAFKA_CFG_PROCESS_ROLES
value: "broker,controller"
- name: KAFKA_CFG_LISTENERS
value: "PLAINTEXT://:9092,CONTROLLER://:9093"
- name: KAFKA_CFG_ADVERTISED_LISTENERS
value: "PLAINTEXT://kafka:9092"
- name: KAFKA_CFG_CONTROLLER_QUORUM_VOTERS
value: "0@kafka:9093"
- name: KAFKA_CFG_CONTROLLER_LISTENER_NAMES
value: "CONTROLLER"
- name: ALLOW_PLAINTEXT_LISTENER
value: "yes"
连接语法
Bootstrap Servers
kafka:9092
Spring Kafka
spring:
kafka:
bootstrap-servers: kafka:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
consumer:
group-id: order-service
auto-offset-reset: earliest
常用 CLI 语法(终端内)
# 创建 topic(3 分区,副本 1 — 单节点)
kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic order-events --partitions 3 --replication-factor 1
# 生产消息
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic order-events
> {"orderId":1001,"status":"paid"}
# 消费
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic order-events --from-beginning
# 查看消费组 lag
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-service
核心配置项
| 参数 | 建议 |
|---|---|
num.partitions | 按并行度,常 3~12 |
retention.ms | 日志保留,默认 7 天可调 |
replication.factor | 单节点仅能为 1;生产多 broker 设 3 |
小紫运维
磁盘占满 → 存储 扩容 + 调 retention;消费 lag 高 → 增加消费者实例(同 group);Bitnami 镜像路径以容器内 ls /opt/bitnami/kafka/bin 为准。
4.3 RocketMQ 5.2(Apache)
适用场景
| 选用 | 典型业务 |
|---|---|
| Java / 阿里系微服务 消息中间件 | 订单、支付、库存异步通知 |
| 事务消息、顺序消息、延迟消息 | 电商下单与扣库存一致性 |
| 与国内团队技术栈一致 | Spring Cloud Alibaba、Dubbo 常用 |
与 Kafka 对比:RocketMQ 低延迟、事务消息 更成熟;日志超大数据流 仍优先 Kafka。与 RabbitMQ 对比:RocketMQ 吞吐更高,管理控制台需另装(商店另有 RocketMQ Dashboard 类条目时可搭配)。
架构说明(必读)
RocketMQ 至少两个角色:
| 角色 | 作用 | 默认端口 |
|---|---|---|
| NameServer | 路由注册、Broker 发现 | 9876 |
| Broker | 消息存储与投递 | 10911(主)、10909(VIP 通道) |
应用商店目录 仅一条 rocketmq 模板(镜像 apache/rocketmq:5.2.0,预填端口 9876)。完整可用集群需部署两次:一次 NameServer,一次 Broker(同一镜像,启动命令不同)。
商店部署参数
| 项 | 值 |
|---|---|
| id | rocketmq |
| 镜像 | apache/rocketmq:5.2.0(基础系统 Debian) |
| 目录预填端口 | 9876(NameServer) |
| PVC | 20 Gi → /home/rocketmq/store |
| 资源 | 中型 · 2 核 4Gi |
部署一:NameServer
- 应用商店 → 中间件 → RocketMQ → 一键部署
- 命名空间
middleware,应用名称rocketmq-namesrv(即 Service 名) - 步骤 ②:启用 PVC
20 Gi,longhorn,挂载/home/rocketmq/store - 步骤 ③:容器端口
9876,Service ClusterIP(或 NodePort30976供集群外调试) - 步骤 ④ 确认页 YAML,在
containers[0]增加启动命令:
command: ["sh", "mqnamesrv"]
- 确认部署 → 就绪
1/1
验收(终端)
# 查看 9876 监听
netstat -tlnp | grep 9876
# 或
ss -tlnp | grep 9876
部署二:Broker
- 再次从商店部署 RocketMQ(同一镜像)
- 应用名称
rocketmq-broker(勿与 namesrv 重名) - PVC
20 Gi(Broker 与 NameServer 各用各的 PVC) - 步骤 ③:容器端口先填
10911,Service ClusterIP - 步骤 ④ YAML:
command: ["sh", "mqbroker"]
args:
- "-n"
- "rocketmq-namesrv:9876"
- "autoCreateTopicEnable=true"
env:
- name: JAVA_OPT_EXT
value: "-Xms512m -Xmx512m -Xmn256m"
rocketmq-namesrv为同 Namespace 下 NameServer 的 Service 名;跨 Namespace 用 FQDN:rocketmq-namesrv.middleware.svc.cluster.local:9876
- 确认部署 后,应用部署 → 端口 补映射(若需集群外或跨节点访问):