跳至主要内容

使用kafka作为消息队列,有哪些容易踩的坑

使用 Kafka 作为消息队列时,有些坑是因为它设计初衷是"分布式日志系统"而非传统消息队列(如 RabbitMQ),如果照搬传统 MQ 的思维方式,很容易踩坑。以下是常见的几类问题:

1. 消息顺序问题

  • 分区内有序,分区间无序:Kafka 只保证同一分区内的消息顺序,如果你的业务需要严格顺序(比如同一用户的操作按顺序处理),必须保证相关消息落到同一分区(合理设计 key)。
  • 重试导致乱序:生产者开启重试(retries)且 max.in.flight.requests.per.connection > 1 时,重试的消息可能"插队",导致同一分区内乱序。如果对顺序敏感,建议设置为 1,或使用 enable.idempotence=true(幂等生产者会自动处理这个问题)。

2. 重复消费 / 消息丢失

  • at-least-once 是默认语义:消费者处理完消息后才提交 offset,如果提交前挂掉,重启后会重复消费。业务逻辑必须做到幂等,否则会出问题(比如重复扣款)。
  • 自动提交的陷阱enable.auto.commit=true 时,offset 是按时间间隔自动提交的,可能在消息还没处理完就提交了(消息丢失),也可能处理完了还没提交就挂了(重复消费)。生产环境通常建议手动提交,并且在业务处理成功后再提交。
  • acks 设置不当acks=0acks=1 在某些故障场景下会丢消息,需要金融级可靠性的场景建议 acks=all 配合 min.insync.replicas

3. 消费者组与 Rebalance 风暴

  • Rebalance 影响吞吐:消费者数量变化、消费者处理超时(超过 max.poll.interval.ms)都会触发 rebalance,期间整个消费者组会短暂停止消费。如果单条消息处理耗时很长,容易频繁触发 rebalance。
  • 消费者数量 > 分区数:多余的消费者会一直空闲,分区数决定了消费并行度的上限,扩容消费者之前要先看分区数够不够。

4. 把 Kafka 当成传统 MQ 用

  • 没有真正的"消息确认/拒绝/重新入队"机制:不像 RabbitMQ 那样可以针对单条消息 ack/nack,Kafka 的消费进度是通过 offset 位移推进的,想做"这条失败了重试,其他继续"就得自己实现(比如死信队列、单独重试 topic)。
  • 没有优先级队列、延迟队列的原生支持:这些都得自己在应用层实现(比如用时间戳 + 定时扫描,或额外引入延迟消息组件)。

5. 分区数设置一旦"想少了"很麻烦

  • 分区数只能增加不能减少,而增加分区会破坏原有的 key-分区映射关系(导致同 key 消息路由到不同分区,破坏顺序性)。初期要预估好吞吐和并行度需求。

6. 消息堆积与磁盘问题

  • 消费跟不上生产速度会导致 lag 持续增长,如果没有监控 lag 指标,问题往往是等到出大事才发现。
  • 日志保留策略(retention):默认按时间或大小滚动删除,如果消费者故障时间过长,数据可能已经被删除,造成永久丢失(不像传统 MQ 消息消费后才删除)。

7. 生产端批量与延迟的权衡

  • linger.msbatch.size 设置不当,要么延迟高,要么吞吐低,需要根据业务场景(实时性 vs 吞吐量)调优。

8. 大消息问题

  • Kafka 默认单条消息大小限制在 1MB 左右(message.max.bytes),如果业务经常产生大消息(比如图片、大 JSON),需要提前调整配置或改用"消息体存对象存储、Kafka 只传引用"的模式。

总结一句话:Kafka 本质是个分布式提交日志(commit log),用它做消息队列时,很多"队列该有的功能"(单条确认、优先级、延迟)都得自己在应用层补上,同时顺序性、幂等性、offset 管理这几块是最容易踩坑、也最需要在设计初期就想清楚的地方。

此博客中的热门博文

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 等机制)。 一句话记忆: 要么全做完、前后不破坏规则、互不干扰、做完不丢。