高级 #distributed#kafka#message-queue

消息队列(Kafka)

消息队列(MQ)解耦生产者和消费者——生产者发消息到队列,消费者异步处理。Kafka 是高性能分布式消息系统,用分区和日志实现高吞吐

📬 发了消息不用等——就像发短信

你给朋友发微信问”今晚吃什么?“——你不会拿着手机等他回复才做下一件事。你发了消息就继续干别的事——他回了他回,你看到了再说。

这就是消息队列的核心思想:发消息的不等收消息的——解耦(Decoupling)。

在分布式系统中,服务 A 需要通知服务 B 做某件事——但不需要 A 等 B 做完。A 把消息发到”队列”里,B 从队列里取消息处理——如果 B 忙,消息在队列里等着——不会把 A 堵住。

🏪 类比:奶茶店的取餐流程

同步方式:你点单→站在柜台前等→奶茶做好了递给你→你走。(你被”阻塞”了——啥也做不了)

异步方式:你点单→拿号→找个座位坐下玩手机→叫到号去取。(你被”解耦”了——可以做自己的事)

消息队列 = 那个”号牌系统”——生产者(你)和消费者(奶茶店)不必同时在线。


📋 消息队列的核心作用

作用说明类比
解耦生产者和消费者不直接依赖发短信而不是打电话
削峰填谷突发流量先存队列,慢慢处理双十一订单先放队列,慢慢发货
异步请求不必等所有处理完成下单后先告诉你”成功”,后台慢慢处理
广播一条消息通知多个消费者班级群发通知——所有人都能收到

🏛️ Kafka——高性能分布式消息系统

Kafka 的核心概念

Topic(主题)= 消息的分类(如"订单消息"、"日志消息")

  ├── Partition 0(分区 0)= 分片,支持并行
  │     └── [msg0, msg1, msg2, ...] ← 有序追加日志

  └── Partition 1(分区 1)
        └── [msg0, msg1, ...]
# Kafka 的简化使用
from kafka import KafkaProducer, KafkaConsumer

# 生产者——发消息
producer = KafkaProducer(bootstrap_servers='localhost:9092')
producer.send('orders', b'{"user": "张三", "amount": 99}')

# 消费者——收消息(从自己的进度继续读)
consumer = KafkaConsumer('orders', bootstrap_servers='localhost:9092')
for msg in consumer:
    print(f"处理订单: {msg.value}")
    process_order(msg.value)

Kafka 的高吞吐秘密

1. 顺序写磁盘——Kafka 写入是日志追加(比随机写快 6000 倍)
2. 分区并行——一个 topic 多个 partition 可同时读写
3. 零拷贝——从磁盘到网卡直接传输,不经用户态内存
4. 批量——消息攒一批再发/再写

📝 小结

概念一句话
消息队列异步解耦——发消息的不用等收消息的
Topic(主题)消息的分类
Partition(分区)分片存储——支持并行读写
Kafka高性能分布式消息系统——顺序写+零拷贝
削峰填谷瞬时高并发→队列缓冲→平滑消费
消费者组组内各消费者分摊不同分区

为什么先学这个? 消息队列是微服务通信的基础设施——微服务架构