中间件:消息队列与 Elasticsearch¶
一句话:消息队列负责"异步解耦、削峰填谷、最终一致",Elasticsearch 负责"全文检索与向量召回",两者撑起了高并发系统的异步与搜索两条主轴。
概念¶
中间件是分布式系统的"连接组织"。本节聚焦两类高频中间件:消息队列(Kafka / RocketMQ) 和 搜索引擎(Elasticsearch)。
消息队列的核心价值是解耦、异步、削峰:生产者把消息扔进队列即可返回,消费者按自己的节奏处理,把瞬时洪峰削平成平滑消费。它也是实现分布式事务最终一致性(本地消息表、事务消息)的载体。Kafka 与 RocketMQ 在设计哲学上差异明显,选型取决于是否需要事务消息、顺序消息与高吞吐。
Elasticsearch 在 AI 时代承担双重角色:传统倒排索引全文检索,以及向量检索(配合 RAG)。理解它的底层结构(倒排索引、分片、segment)是优化召回与延迟的基础。
原理¶
Kafka vs RocketMQ¶
| 维度 | Kafka | RocketMQ |
|---|---|---|
| 出身 | LinkedIn,高吞吐日志/流 | 阿里,脱胎自电商交易 |
| 吞吐 | 极高(百万级/单机,顺序写+零拷贝) | 高,略低于 Kafka |
| 延迟 | 批处理优化,毫秒级 | 更低,支持近实时 |
| 事务消息 | 有(精确一次,生产者事务) | 原生支持(半消息+回查,业务级事务) |
| 顺序消息 | 分区内有序 | 支持全局/分区顺序,电商场景成熟 |
| 延时消息 | 需自建(特定 level) | 原生支持多级延时 |
| 消息回溯 | 按 offset/时间戳 | 支持按时间回溯 |
| 适用 | 日志、大数据流、事件溯源 | 交易、订单、金融、事务消息场景 |
核心区别总结:Kafka 为"高吞吐流式日志"而生,RocketMQ 为"电商交易、事务消息、延时消息"而生。需要业务级事务消息、可靠顺序消费选 RocketMQ;纯高吞吐数据管道选 Kafka。
RocketMQ 事务消息(两阶段)¶
RocketMQ 的事务消息解决"本地事务与消息发送的原子性",流程是两阶段:
- 发送半消息(Half Message):生产者先发一条对消费者不可见的半消息到 broker。
- 执行本地事务:半消息发送成功后,执行本地 DB 事务。
- 提交/回滚:本地事务成功 → 提交消息(对消费者可见);失败 → 回滚消息。
- 回查(事务回查):若生产者宕机未提交/回滚,broker 主动回查生产者"这个事务到底成了没",生产者根据本地事务状态答复。
这与本地消息表方案对比:本地消息表是"业务表+消息表同库写入 → 定时扫描消息表发送",强依赖定时轮询,一致性依赖本地事务。事务消息省了消息表与轮询,靠 broker 回查驱动,更轻量但要求实现 TransactionListener。
消息可靠性与顺序¶
可靠性三段保证:
| 环节 | 机制 |
|---|---|
| 生产端 | 同步发送 + 重试 + ACK 确认(producer 收到 broker 确认才算成功) |
| Broker | 持久化到磁盘 + 副本/主从同步复制 |
| 消费端 | 手动 ACK(消费成功才确认 offset),失败重试+死信队列 |
顺序消息:只有"同一业务 key 的消息发到同一队列、且单队列单消费者消费"才能保证顺序。代价是牺牲并行度。全局顺序需单队列,吞吐极低,实际多用分区顺序(如按 orderId 分区)。
Elasticsearch 倒排索引¶
ES(基于 Lucene)的核心是倒排索引:
- 正排索引:文档 ID → 文档内容。
- 倒排索引:词条 (Term) → 包含该词的文档 ID 列表(Postings List)。
分词器把文档切成词条,每个词条记录出现的文档 ID 与位置/频率。查询时先查倒排索引拿到候选文档集,再按相关度(BM25/TF-IDF)打分排序。倒排索引让全文检索从"逐篇扫描"变成"字典查询",复杂度骤降。
ES 的其他关键结构:分片 (Shard)(水平扩展单元,每个分片是一个独立 Lucene 索引)、segment(不可变的倒排索引段,合并优化)、FST(前缀字典,压缩 term 查找)。
向量检索 (kNN)¶
ES 从 7.x 起支持稠密向量字段(dense_vector)与近似最近邻 (ANN) 检索:
- 存文档的 embedding 向量。
- 用 HNSW 算法(8.x 起)做 ANN 检索,查询时返回与查询向量最近似的 K 个文档。
- 配合 RAG,可把向量召回与传统倒排召回用 RRF (Reciprocal Rank Fusion) 融合,提升召回质量(报告项目一的混合检索正是此模式)。
实战要点¶
结合报告"Kafka、RocketMQ(事务消息)、Elasticsearch(含向量检索)"技能与项目实战:
-
选型看场景:电商交易/订单/支付这类需要事务消息、可靠顺序、延时消息的业务用 RocketMQ;日志采集、行为埋点、大数据流这种追求极致吞吐的用 Kafka。混用是常见的工程现实。
-
消息必须手动 ACK,幂等消费:自动 ACK 在消费失败时会丢消息。用手动 ACK + 失败重试 + 死信队列兜底。同时消费者必须幂等(用唯一键/状态机去重),因为网络重试会导致重复消费。
-
顺序消息只在必要时用:顺序消费会大幅降低吞吐。先问业务是否真的需要严格顺序(如"先创建订单再扣库存"),能用"事件时间戳排序"或"业务幂等"解决的,就不要上顺序消息。
-
事务消息要实现好回查:RocketMQ 事务消息的回查接口必须能根据业务主键查到本地事务的最终状态(成功/失败/未知),且回查要幂等。回查接口设计不好会导致半消息悬挂。
-
积压治理要分级:消息积压时,盲目加消费者可能打爆下游。要按下游承受能力限速消费、或临时扩容分区、或用"快速消费+落库后异步处理"先消化队列。详见深度题"积压百万如何优雅消费不击穿下游"。
-
ES 大结果集别用深分页:
from + size超过 10000 会触发性能问题(要查所有分片的前 N 条归并)。改用search_after(游标)或scroll(遍历)。向量检索用num_candidates控制 ANN 召回范围平衡精度与延迟。
本节相关题目¶
| 难度 | 题目 | 链接 |
|---|---|---|
| 基础 | Kafka vs RocketMQ 核心区别 | → 题库 |
| 进阶 | RocketMQ 事务消息两阶段 vs 本地消息表 | → 题库 |
| 深度 | 消息队列积压百万如何优雅消费不击穿下游 | → 题库 |