学后端时,MQ 和 Kafka 是两个非常容易混淆的东西。
刚开始看,会发现它们长得几乎一样:
Producer
↓
消息中间件
↓
Consumer
都有生产者,都有消费者,都能发送消息。
于是很容易产生几个疑问:
Kafka 和 MQ 不都是生产消息、消费消息吗?
Kafka 消费以后还能重新消费,那是不是比 MQ 更强?
既然 Kafka 什么都能做,那为什么还要 MQ?
真正理解 Kafka 和 MQ,不能从 Producer、Consumer、Broker 这些名词开始,而应该先理解它们分别想解决什么问题。
可以先记住两个最简单的概念:
MQ 更像任务队列。
Kafka 更像分布式系统中的事件流水。
这两个概念一旦建立起来,后面的 Queue、Topic、Partition、Offset、Consumer Group 就会容易很多。
一、先理解最普通的 MQ
假设系统需要生成一份机器人清扫报告。
用户在 App 中点击:
生成清扫报告
如果生成 PDF 要花 5 秒钟,接口完全可以同步执行:
App
↓
Spring Boot
↓
生成 PDF
↓
5 秒以后
↓
返回 App
但这样用户就得一直等待。
更常见的做法是:
App
↓
Spring Boot
↓
MQ
↓
报告服务
↓
生成 PDF
Spring Boot 只是往 MQ 中扔一个任务:
给机器人 001 生成清扫报告
然后直接告诉 App:
任务已提交
后台的报告服务再慢慢处理。
因此这里 MQ 中保存的东西,本质上是:
等待某个消费者完成的任务。
可以暂时把 MQ 理解成:
MQ
=
任务队列
二、MQ 没有消费者会怎么样?
假设 Producer 连续发送三条消息:
Producer
↓
A
B
C
↓
Queue
如果这时候没有 Consumer,消息并不会凭空消失。
而是:
Queue
[A] [B] [C]
暂时排在那里。
过一会消费者启动:
Consumer
↓
取出 A
↓
处理成功
↓
ACK
然后:
取出 B
取出 C
直到任务全部完成。
因此 MQ 很像一个:
任务待办池。
没人处理的时候:
任务先排队
有人处理的时候:
取任务
↓
执行
↓
成功
↓
ACK
所以 MQ 里的消息并不是必须立刻消费。
恰恰相反:
“先存着,等消费者有能力的时候再处理”,本身就是 MQ 存在的重要意义。
三、为什么业务系统里经常需要 MQ?
例如用户下单以后,可能需要同时做很多事情:
创建订单
发送短信
发送邮件
增加积分
发送优惠券
通知仓库
如果全部同步执行:
用户
↓
订单服务
↓
短信服务
↓
邮件服务
↓
积分服务
↓
仓库服务
↓
返回
整个接口会越来越慢。
于是可以变成:
用户下单
↓
订单保存成功
↓
发送 MQ 消息
↓
立即返回
后台:
MQ
↓
┌───────┼────────┐
↓ ↓ ↓
短信服务 积分服务 其他服务
这样订单服务就不需要等待所有业务执行完成。
这就是 MQ 一个非常典型的作用:
异步
把不需要立即完成的事情放到后台执行。
除此之外,MQ 还有两个非常重要的作用。
四、MQ 可以用来解耦
假设订单服务直接调用短信服务:
订单服务
↓
短信服务
订单服务必须知道:
短信服务在哪里
短信接口是什么
短信服务有没有挂
以后又增加邮件:
订单服务
├→ 短信服务
└→ 邮件服务
再增加积分:
订单服务
├→ 短信
├→ 邮件
└→ 积分
订单服务和越来越多系统发生依赖。
如果改成:
订单服务
↓
MQ
订单服务只需要告诉系统:
订单创建成功
至于谁处理,由消费者自己决定。
这样系统之间的依赖就降低了。
这就是:
解耦。
五、MQ 还可以削峰
假设 MySQL 最多能够稳定处理:
1000 个请求 / 秒
突然秒杀活动来了:
10000 个请求 / 秒
如果所有请求直接进入 MySQL:
10000 请求
↓
MySQL
数据库很可能被瞬间打垮。
这时候可以在中间增加 MQ:
10000 请求
↓
MQ
↓
消费者按照系统能力处理
↓
1000 / 秒
↓
MySQL
MQ 就像一个水库。
洪峰先进入水库:
10000
再按照下游能够接受的速度慢慢放出去:
1000
1000
1000
…
这就是:
削峰填谷。
所以业务系统里 MQ 的使用场景其实非常多。
常见的有:
发送短信
发送邮件
生成报表
图片处理
异步计算
延迟任务
订单处理
失败重试
削峰
这些场景有一个共同特点:
系统希望某件事情最终被执行完成。
六、Kafka 解决的问题开始不一样了
再来看机器人系统。
机器人运行过程中会不断产生数据:
10:00:01 位置 A 电量 80%
10:00:02 位置 B 电量 80%
10:00:03 位置 C 电量 79%
10:00:04 位置 D 电量 79%
机器人还可能不断产生:
位置变化
速度变化
任务状态
电量变化
设备上线
设备掉线
故障
运行日志
传感器数据
这些东西和:
生成一份 PDF
完全不是一个性质。
生成 PDF 是:
请别人帮我完成一个任务。
机器人位置变化则是:
系统刚刚发生了一件事情。
这就是 Kafka 开始擅长的领域:
事件。
七、Kafka 可以理解成一份不断追加的事件流水
假设机器人不断产生事件:
offset 0 RobotConnected
offset 1 PositionChanged
offset 2 BatteryChanged
offset 3 PositionChanged
offset 4 RobotError
offset 5 PositionChanged
Kafka 会把这些事件不断追加进去。
可以把它想象成一个巨大的日志:
0
1
2
3
4
5
6
7
8
…
消费者并不是简单地把里面的数据“拿走”。
而是在读取:
我现在读到哪里了?
例如状态服务:
已经读到 offset 5
数据分析服务:
读到 offset 3
日志服务:
读到 offset 5
它们各自拥有自己的消费进度。
这就是 Kafka 和传统任务队列非常不一样的地方。
八、为什么 Kafka 消费完还能重新消费?
传统 MQ 可以先粗略理解成:
消息
↓
Queue
↓
Consumer
↓
处理成功
↓
ACK
↓
任务结束
MQ 更关注:
这个任务有没有被处理成功?
Kafka 则更像:
事件
↓
写入 Kafka
↓
保存在 Partition 中
↓
Consumer 读取
↓
记录 offset
消费者读完以后,并不意味着这条数据立即被删除。
Kafka 中的数据一般按照:
时间
容量
保留策略
等规则保存。
因此原来消费者:
offset = 100
完全可以重新调整到:
offset = 50
然后:
50
51
52
53
…
100
重新读取。
所以 Kafka 的消费更像:
读日志。
而不是:
把任务取走。
九、Kafka 为什么特别适合分布式、多进程系统?
这其实是理解 Kafka 非常重要的一步。
现代后端往往不是只有一个进程。
例如机器人云平台可能有:
设备状态服务
告警服务
日志服务
任务服务
数据分析服务
轨迹服务
这些都是:
不同服务
不同进程
甚至部署在不同服务器
现在机器人发生了一件事情:
Robot001
电量降低到 20%
很多服务都想知道:
状态服务:
我要更新当前电量
告警服务:
我要判断是否低电量报警
日志服务:
我要记录
数据分析服务:
我要拿去统计
如果机器人分别调用四个系统:
Robot
├→ 状态服务
├→ 告警服务
├→ 日志服务
└→ 分析服务
系统耦合会非常严重。
于是可以变成:
Robot001
BatteryChanged(20%)
↓
Kafka
↓
┌──────┼──────┬──────┐
↓ ↓ ↓ ↓
状态 告警 日志 数据分析
服务 服务 服务 服务
机器人只产生一次事件。
Kafka 把这份事件流提供给多个独立消费者。
每个消费者:
按照自己的速度读取
并且:
互相不影响
因此 Kafka 非常适合:
分布式系统中,让多个独立服务/进程消费同一份事件数据流。
这个理解比简单说:
Kafka 是消息队列。
要更加接近 Kafka 的实际价值。
十、但 Kafka 不是“共享变量”
这里有一个容易继续混淆的地方。
既然多个进程都能从 Kafka 获取数据,那么是不是:
Kafka 就是分布式共享数据?
不能完全这么理解。
假设多个服务想知道:
Robot001 当前是否在线?
或者:
用户当前登录 Session 是什么?
或者:
某个 Token 是否有效?
这些关注的是:
当前值。
这种情况下:
Service A ─┐
Service B ─┼──→ Redis
Service C ─┘
Redis 往往更合适。
Kafka 更关注:
Robot001 什么时候上线了?
什么时候掉线了?
什么时候再次上线了?
它保存的是:
事情发生的过程。
因此可以形成一个非常重要的区分。
十一、Redis 和 Kafka 的区别:现在是什么 vs 发生过什么
例如:
Robot001 当前电量:20%
这是一个:
当前状态
可以放 Redis。
而:
10:00 电量 30%
10:10 电量 25%
10:20 电量 20%
这是:
状态变化事件流
可以进入 Kafka。
所以可以简单记成:
Redis
=
现在是什么?
而:
Kafka
=
发生过什么?
这两个东西解决的问题其实完全不同。
十二、再把 MySQL 放进来
这时候就能把几个后端常见组件放到一起理解了。
假设机器人系统中存在这些数据。
机器人基本资料
机器人名称
SN
设备型号
所属用户
创建时间
这些是正式业务数据:
→ MySQL
机器人当前状态
在线
当前电量
当前任务
当前坐标
这些经常读取,而且变化频繁:
→ Redis
机器人运行事件
上线
掉线
位置变化
电量变化
故障
日志
传感器数据
这些形成持续的数据流:
→ Kafka
后台任务
生成报告
发送短信
发送邮件
图片处理
这些是等待执行的任务:
→ MQ
于是系统开始变得非常清楚。
十三、完整来看四种组件
可以形成这样的第一层架构直觉:
分布式系统
正式业务数据
↓
MySQL
当前共享状态
↓
Redis
异步业务任务
↓
RabbitMQ / RocketMQ
大量事件数据流
↓
Kafka
可以进一步记成四句话:
MySQL:
业务最终是什么?
Redis:
现在是什么?
MQ:
有什么事情需要做?
Kafka:
刚刚发生了什么?
这四句话非常适合建立后端系统的第一层架构认知。
十四、机器人系统中的完整例子
假设机器人现在电量从:
21%
变成:
20%
机器人上报:
BatteryChanged
robotId = 001
battery = 20
事件进入 Kafka:
Robot
↓
Kafka
然后多个系统消费:
Kafka
↓
┌─────────┼──────────┐
↓ ↓ ↓
状态服务 告警服务 数据分析
状态服务:
更新 Redis:
robot001:battery = 20
告警服务发现:
battery <= 20%
于是产生一个任务:
发送低电量通知
这个任务进入 MQ:
告警服务
↓
MQ
↓
通知服务
↓
发送 Push / 短信
同时正式的告警记录:
写入 MySQL
最终整个链路就是:
Robot
↓
BatteryChanged
↓
Kafka
↓
告警服务
↓
发现低电量
↓
MQ
↓
通知服务
同时:
状态服务 → Redis
告警记录 → MySQL
到这里可以发现:
Kafka、MQ、Redis、MySQL 根本不是竞争关系。
它们分别负责不同的问题。
十五、MQ 和 Kafka 最核心的区别
现在回过头看,可以非常简单地总结。
MQ
更关注:
事情有没有被处理完成?
例如:
发送短信
发送邮件
生成报告
图片转换
后台计算
因此可以粗略理解为:
Task
↓
MQ
↓
Consumer
↓
完成任务
Kafka
更关注:
系统发生了什么?哪些服务需要获取这些事件?
例如:
设备上线
设备掉线
机器人位置变化
用户点击
订单状态变化
系统日志
传感器数据
因此可以理解为:
Event
↓
Kafka
↓
多个独立服务
↓
按照自己的进度消费
十六、Kafka 并不是“更高级的 MQ”
这是非常容易出现的误区。
看到 Kafka:
吞吐量高
可以持久化
可以重复消费
可以多个消费者读取
很容易产生:
既然 Kafka 什么都能做,为什么还需要 MQ?
事实上正确的问题并不是:
Kafka 和 MQ 谁更厉害?
而应该是:
我现在解决的是任务处理问题,
还是事件流问题?
如果是:
帮我完成一件事情。
优先想到:
MQ
如果是:
系统发生了一件事情,我需要让多个分布式服务知道。
而且:
数据量大
持续产生
需要保留
可能重复读取
多个消费者独立消费
那么 Kafka 就非常合适。
当然实际系统并不是绝对的,Kafka 也可以承担某些异步任务,RabbitMQ、RocketMQ 也支持发布订阅。
但作为学习阶段的第一层架构判断:
MQ = 任务队列
Kafka = 分布式事件流水
这个模型非常实用。
十七、最后形成一张完整的架构图
以机器人云平台为例:
机器人云平台
Robot
↓
位置 / 电量 / 状态 / 日志 / 故障
↓
Kafka
↓
┌───────────┬────────────┬─────────────┐
↓ ↓ ↓ ↓
状态服务 告警服务 日志服务 数据分析
↓
Redis
保存当前状态
告警服务
↓
发现需要通知
↓
MQ
↓
通知服务
↓
Push / 短信 / 邮件
设备资料 / 用户 / 任务 / 告警记录
↓
MySQL
这时候几个组件的位置就非常清楚:
Redis
→ 当前共享状态
MySQL
→ 正式业务数据
MQ
→ 等待执行的任务
Kafka
→ 分布式服务共享的事件数据流
十八、总结
最开始学习 MQ 和 Kafka 时,很容易陷入:
Producer
Consumer
Queue
Topic
Broker
Partition
Offset
Consumer Group
ACK
这些技术名词。
但如果先从系统设计的角度理解,会简单很多。
遇到一个需求时,可以先问:
这是需要别人完成的一件事情吗?
例如:
发短信
生成报告
处理文件
考虑:
MQ
这是系统中发生的一件事情吗?
例如:
机器人位置变化
设备掉线
电量变化
系统日志
并且多个分布式服务都需要消费:
考虑:
Kafka
我需要共享的是当前状态吗?
例如:
机器人当前在线状态
当前电量
Session
Token
考虑:
Redis
这是最终需要长期保存的业务数据吗?
例如:
用户
设备
订单
任务
告警记录
考虑:
MySQL
最后可以浓缩成四句话:
MySQL:业务数据最终是什么?
Redis:现在是什么?
MQ:有什么事情需要做?
Kafka:刚刚发生了什么,并且哪些分布式服务需要知道?
理解到这一层以后,再学习 Broker、Partition、Offset、Consumer Group,就不再只是背概念了。
因为你已经知道:
这些机制最终都是为了让分布式系统中的任务和事件可靠地流动起来。


