跳转至

中间件:消息队列与 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 的事务消息解决"本地事务与消息发送的原子性",流程是两阶段:

  1. 发送半消息(Half Message):生产者先发一条对消费者不可见的半消息到 broker。
  2. 执行本地事务:半消息发送成功后,执行本地 DB 事务。
  3. 提交/回滚:本地事务成功 → 提交消息(对消费者可见);失败 → 回滚消息。
  4. 回查(事务回查):若生产者宕机未提交/回滚,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(含向量检索)"技能与项目实战:

  1. 选型看场景:电商交易/订单/支付这类需要事务消息、可靠顺序、延时消息的业务用 RocketMQ;日志采集、行为埋点、大数据流这种追求极致吞吐的用 Kafka。混用是常见的工程现实。

  2. 消息必须手动 ACK,幂等消费:自动 ACK 在消费失败时会丢消息。用手动 ACK + 失败重试 + 死信队列兜底。同时消费者必须幂等(用唯一键/状态机去重),因为网络重试会导致重复消费。

  3. 顺序消息只在必要时用:顺序消费会大幅降低吞吐。先问业务是否真的需要严格顺序(如"先创建订单再扣库存"),能用"事件时间戳排序"或"业务幂等"解决的,就不要上顺序消息。

  4. 事务消息要实现好回查:RocketMQ 事务消息的回查接口必须能根据业务主键查到本地事务的最终状态(成功/失败/未知),且回查要幂等。回查接口设计不好会导致半消息悬挂。

  5. 积压治理要分级:消息积压时,盲目加消费者可能打爆下游。要按下游承受能力限速消费、或临时扩容分区、或用"快速消费+落库后异步处理"先消化队列。详见深度题"积压百万如何优雅消费不击穿下游"。

  6. ES 大结果集别用深分页from + size 超过 10000 会触发性能问题(要查所有分片的前 N 条归并)。改用 search_after(游标)或 scroll(遍历)。向量检索用 num_candidates 控制 ANN 召回范围平衡精度与延迟。

本节相关题目

难度 题目 链接
基础 Kafka vs RocketMQ 核心区别 → 题库
进阶 RocketMQ 事务消息两阶段 vs 本地消息表 → 题库
深度 消息队列积压百万如何优雅消费不击穿下游 → 题库