跳至主要内容

Kafka架构和消息流转详解

一、Kafka是什么

一句话:Kafka是一个分布式的、基于日志(log)的消息系统。它的核心思想很简单——把消息像写日志一样"追加"到磁盘文件里,消费者按顺序"读"这个文件,仅此而已。这个简单模型是它高吞吐的根本原因。

二、核心架构组件

在讲数据流之前,先认识几个名词:

  • Producer(生产者):发消息的一方,比如你的业务代码
  • Broker:Kafka的服务器节点,多个Broker组成一个Cluster(集群)
  • Topic(主题):消息的分类,比如"订单创建"topic
  • Partition(分区):每个Topic会被切成多个分区,分区是Kafka并行处理和存储的最小单位
  • Replica(副本):每个分区有多个副本分布在不同Broker上,其中一个是Leader(负责读写),其余是Follower(只同步数据,不对外服务)
  • Consumer(消费者) / Consumer Group(消费组):消费消息的一方,同一个组内的消费者共同"瓜分"一个Topic的所有分区
  • Controller:集群里选出的一个特殊Broker,负责管理分区Leader选举等元数据工作(以前靠ZooKeeper协调,新版本Kafka用内置的KRaft协议取代了ZooKeeper)## 三、数据从Producer到Broker再到Consumer的完整过程

  • Producer发送消息:业务代码调用send(),消息先进入客户端本地的一个内存缓冲区(RecordAccumulator),并不是发一条就立刻走网络

  • 选择分区:如果你指定了key,Kafka默认对key做hash来决定发到哪个分区(保证同key的消息永远进同一个分区,从而保证顺序);没指定key则轮询或用粘性分区策略
  • 批量发送:客户端把同一分区的多条消息打包成一个batch,达到大小或时间阈值后一起发给对应分区的Leader Broker(这是Kafka高吞吐的关键之一——减少网络往返)
  • Leader写入日志:Leader收到消息后,以顺序追加的方式写入本地磁盘的log文件,并分配一个递增的offset
  • 副本同步:Follower副本主动从Leader拉取新消息进行同步,当所有处于ISR(In-Sync Replicas,同步副本集合)中的副本都写完后,才算这条消息"提交"(commit)
  • 返回ack:根据Producer的acks配置(0/1/all)决定什么时候返回成功,这一步直接决定了可靠性和性能的取舍
  • Consumer拉取:消费者是主动poll(拉模式,不是broker推给你),按上次记录的offset位置继续往后读
  • 提交offset:消费完成后,消费者把"我读到哪了"这个offset提交到一个特殊的内部Topic __consumer_offsets里,下次重启就从这个位置接着消费

四、日志和数据结构细节

Kafka的性能秘密基本都藏在这套存储结构里:

  • Partition = 一个只追加的文件序列:每个分区在磁盘上对应一个目录,消息按offset严格递增地写入
  • Log Segment(日志段):一个分区的数据不会塞进一个文件,而是切成多个segment,每个segment有三类文件:
    • .log:真正存消息内容的文件
    • .index:offset到物理文件位置的稀疏索引,用于快速定位
    • .timeindex:时间戳到offset的索引,支持"按时间查消息"
  • 消息本身的结构包含:key、value、时间戳、offset、headers、CRC校验、压缩标记等
  • 顺序写 + OS Page Cache + 零拷贝(sendfile):这三者结合是Kafka比传统消息队列快得多的根本原因——磁盘顺序写接近内存速度,读取时利用系统缓存和零拷贝减少用户态/内核态数据拷贝
  • 数据清理两种策略:按时间/大小滚动删除旧segment(retention),或者Log Compaction(只保留每个key最新的一条,常用于存"最新状态"场景,比如Kafka自己用它存offset)## 五、实际使用中要注意的坑

  • 分区数不是越多越好:分区数决定了并行度上限,但太多分区会增加Controller管理开销、增加端到端延迟、增加文件句柄消耗。一般根据目标吞吐量和消费者数量来定,而不是越多越好

  • 分区键设计要避免热点:如果大量消息用同一个key(比如都用"CN"作为地区key),会导致某个分区数据量远超其他分区,造成热点和消费延迟不均
  • acks参数是可靠性与性能的权衡:acks=0最快但可能丢数据;acks=1只要leader写成功就返回,leader挂了还没同步给follower的数据会丢;acks=all配合min.insync.replicas才能真正保证不丢消息,但延迟更高
  • 顺序性只在分区内保证,跨分区不保证顺序。如果业务要求严格顺序,必须把相关消息发到同一个分区
  • 消费者组的Rebalance是常见坑:新消费者加入、消费者挂掉、session.timeout.ms太短都会触发rebalance,期间会短暂停止消费。处理消息耗时过长导致max.poll.interval.ms超时是最常见的诱因,要么加快处理速度,要么调大这个参数,要么用异步处理+手动提交offset
  • 重复消费是常态,不是异常:Kafka是"至少一次"(at-least-once)语义为主,消费者重启、rebalance都可能导致重复消费,业务逻辑必须做幂等处理(比如用消息里的唯一ID做去重)
  • 要监控消费延迟(Lag):consumer的offset落后producer offset太多说明消费能力跟不上,这是生产环境最需要盯的指标之一
  • 磁盘容量规划:Kafka数据是先落盘再清理的,retention设置太长+topic量大很容易把磁盘写爆
  • 序列化格式建议提前定好:生产环境不建议直接传JSON字符串,用Avro/Protobuf配合Schema Registry能省很多后期兼容性的麻烦

六、延时队列怎么实现

关键点:Kafka本身并不原生支持"延时消息",它只是一个按offset顺序读写的日志,没有"消息到期才可见"这种机制。要在Kafka上做延时队列,常见思路是:

  1. 多级延时Topic + 轮询转发:仿照RocketMQ思路,建几个不同粒度的延时topic(比如delay-5sdelay-30sdelay-1m),生产者把消息连同"目标可见时间"一起写入对应topic,由一个专门的转发服务不断轮询这些topic,发现到期的消息就转发到真正的业务topic;没到期的可以选择sleep等待或者重新写回原topic(类似时间轮的思路)
  2. 借助外部存储做调度器:把消息连同到期时间存入Redis的ZSET(score=到期时间戳)或数据库表,用一个独立的定时任务/时间轮扫描到期的记录,到期后再发到Kafka真正的topic。这是目前生产环境用得最多的方式,因为Kafka本身不适合频繁"回查"
  3. 本质上都是"外挂"方案:因为Kafka缺少消息级别的可见性控制(不像队列型系统可以"暂扣"某条消息),延时能力必须靠上层组件补齐,精度和维护成本都比原生支持延时的MQ要高

七、和RocketMQ、RabbitMQ的区别

维度 Kafka RocketMQ RabbitMQ
核心模型 分布式日志,partition顺序追加 commitlog + 多队列,借鉴Kafka又做了增强 AMQP协议,交换机(Exchange)+队列路由
消费模式 纯拉模式(pull) 推拉结合(长轮询) 推模式为主
延时消息 不原生支持,需外部方案 原生支持,早期版本固定18个延时级别(1s~2h),新版本支持任意精度延时 通过TTL+死信队列实现,或装delayed-message-exchange插件支持任意延时
事务消息 支持(较复杂,主要面向流处理场景) 原生支持,半消息+回查机制,是其强项之一 支持,但性能开销较大
顺序消息 分区内有序 支持全局顺序或分区顺序 单队列内基本有序,但集群/多消费者场景较弱
吞吐量 最高,适合海量日志/流数据 高,略低于Kafka 相对最低,适合中小规模、低延迟场景
路由能力 简单(topic+partition) 中等(tag过滤) 最强,支持direct/topic/fanout/headers多种灵活路由
典型场景 日志采集、流处理、大数据管道 电商订单、金融交易(强调可靠性和延时/事务消息) 微服务间的低延迟RPC、任务分发

简单总结一下选型直觉:要吞吐量和流处理选Kafka;要原生延时消息、事务消息、订单类强可靠性场景选RocketMQ;要灵活路由、复杂场景解耦、开发效率优先选RabbitMQ

此博客中的热门博文

Elasticsearch 读写原理指南

### 1. 什么是 segment,里面装了什么? 在 Lucene(也是 Elasticsearch)里,索引被切分成若干 **segment(段)**,每个 segment 是一个完整的、只读的倒排索引单元。一个 segment 包含: * **倒排词典** —— 用 **FST(Finite‑State Transducer)** 以高度压缩的形式保存每个字段出现的所有 term 以及 term→ord 的映射。对应的磁盘文件是 `*.tim`(新版)或 `*.tis/*.tii`(旧版)。 * **倒排列表(postings)** —— 保存每个 term 出现的文档 ID、频次、位置信息等,文件名通常是 `*.doc`、`*.pos`、`*.pay`。 * **存储字段**(_source、store:true 的字段)—— 以二进制块的形式写入 `*.fdt` / `*.fdx`。 * **doc‑values、norms、向量** 等辅助结构,分别保存在 `*.dv`、`*.norm`、`*.tv` 等文件里。 * **deleted‑docs bitmap**(`*.del`),标记哪些文档已被删除或被更新。 所有这些文件在 segment **写入磁盘后即成为只读**,后续的查询只能读取,永远不会在原文件上进行增删改。 --- ### 2. 原始文档和 FST 为什么都在 segment 里? * **原始文档**:Elasticsearch 默认把完整的 JSON(_source)以及任何 `store:true` 的字段写入 segment 的 `*.fdt/*.fdx` 文件。每个 segment 保存自己的那部分文档,旧的 segment 在合并前仍然保留,直到合并后被删除。 * **FST**:每个字段的词典在每个 segment 中单独维护,采用 FST 进行前缀共享和字节压缩。这样即使同一个 term 在多个 segment 中出现,也会在每个 segment 里拥有独立的映射,查询时只需要在对应 segment 的 FST 中定位即可。 --- ### 3. 查询时到底是怎么遍历 segment 的? 1. **请求入口**      客户端的搜索请求先到达 **协调节点**,协调节点把请求 ...

LLM缓存详解

 可以把“大模型缓存”理解成: 把已经算过的结果(或中间结果)存下来,下次尽量复用 。但这里面其实分几层,不只是简单的“问题→答案”缓存。 1️⃣ 常见的几种缓存类型 (1)KV Cache(推理内部缓存) Transformer 在生成时,会把前面 token 的 Key/Value 向量 缓存下来。 本质:避免重复计算 attention 作用: 同一请求内部加速 特点: 👉 只对“同一上下文继续生成”有效 👉 不跨用户、不跨请求 这类缓存是你体感“流式输出越来越快”的原因之一。 (2)Prompt Cache(提示词缓存) 缓存的是: 相同(或高度相似)的 prompt → 对应的中间表示 / 输出 典型场景: 系统提示词(system prompt)很长 多轮对话里前文基本不变 👉 这里能省掉 前缀计算成本(prefill) (3)Embedding / 语义缓存(Semantic Cache) 这个才是你问题的关键 👇 不是按“字符串完全一致”,而是: 把问题转成向量 → 找“语义相似”的历史问题 → 直接复用答案 2️⃣ 为什么命中缓存成本低很多? 因为大模型推理成本主要在两块: (1)Prefill(吃 prompt) 复杂度 ~ O(n²) 很贵(尤其长 prompt) (2)Decode(逐 token 生成) 每个 token 都要算一遍模型 而缓存命中后: KV cache:不用重复 attention Prompt cache:不用重新 encode 语义缓存: 直接跳过模型推理 👉 相当于从: 几十~几百毫秒 + GPU算力 变成: 一次向量检索(毫秒级)+ 直接返回 所以成本差一个数量级是正常的。 3️⃣ “每个人问法不同,怎么命中缓存?” 这是核心难点,也是工程重点👇 ❌ 不能靠字符串匹配 比如: “今天天气怎么样” “今天外面热不热” 字符串完全不同 → 必须 miss ✅ 用语义相似度(Embedding) 流程一般是: 把问题转 embedding(向量) 在向量数据库里找 TopK 相似问题 如果相似度 > 阈值(比如 0.9) 直接返回缓存答案 一个简单示意 Q1: 北京天气怎么样 → embedding A Q2: 北京今天热吗 → embedding B cosine(A, B) ≈ 0.95...

事务的ACID是什么

 事务的 ACID 是数据库事务必须满足的四个基本性质,用来保证在并发和故障情况下数据的正确性与可靠性: A(Atomicity,原子性) 一个事务中的操作要么 全部成功 ,要么 全部失败回滚 ,不存在“只做了一半”的中间状态。 C(Consistency,一致性) 事务执行前后,数据库都必须处于 一致的合法状态 ,满足约束(如主键、外键、唯一性、业务规则等)。 I(Isolation,隔离性) 并发执行的多个事务之间 相互隔离 ,一个事务未提交的中间结果对其他事务不可见(具体强弱由隔离级别决定)。 D(Durability,持久性) 一旦事务提交成功,其结果会被 永久保存 ,即使系统崩溃也不会丢失(通常依赖 WAL/redo log 等机制)。 一句话记忆: 要么全做完、前后不破坏规则、互不干扰、做完不丢。