企业级知识库项目:基于消息队列的异步 RAG 流水线
上节更了文档解析入库。
这节我们做一下文档发布后的流程,也就是向量化。
我们用消息队列 RabbitMQ 异步来做。

发布后有三个消费者:
Search:收到 INDEX 消息(带整篇元数据 + 正文)→ 写入 Elasticsearch kh_document,供文档级关键词检索
RAG:收到「按文档 ID 重建」消息 → 查库拉正文 → 分块 → Embedding → 写入 Elasticsearch kh_chunk(含向量),供语义检索
KG:收到「按文档 ID 建图」消息 → 查库 → 分块 → 抽实体/关系 → 写入 Neo4j,供知识图谱查询
总之,就是文档发布后,分别做存入向量数据库、全文检索库、图数据库。
异步的用 mq 的消费者来做,三条管线并行、互不影响。
发布接口只负责把 Postgres 里的文档标成已发布,并投递消息
这节先把 RAG 向量化这条链路讲清楚:
也就是:
发布 → 投递 → 消费 → 分块 → 向量化 → 写入 kh_chunk

现在代码仓库迁移到了 gitcode,可以加我微信 guangguangsunlight 开仓库权限,是在 knowledge-hub-backend 下,按照 v1、v2、v3 分支对应每个章节的代码
首先我们要改一下 docker compose 文件,跑一下相关容器:
# RabbitMQ — 文档发布后异步管线(RAG / KG / ES)
rabbitmq:
image: rabbitmq:3.13-management
container_name: knowledge_hub_rabbitmq
restart: unless-stopped
ports:
- "5672:5672" # AMQP
- "15672:15672" # Management UI
environment:
RABBITMQ_DEFAULT_USER: guest
RABBITMQ_DEFAULT_PASS: guest
volumes:
- ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/rabbitmq:/var/lib/rabbitmq
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 10s
timeout: 10s
retries: 5
# Elasticsearch 8.17.0 + IK 中文分词(内置到镜像)
es:
build: ./elasticsearch # 从本地 Dockerfile 构建镜像(自带IK)
container_name: knowledge_hub_elasticsearch
ports:
- "9200:9200" # ES 访问端口
environment:
- discovery.type=single-node # 单节点运行(开发环境)
- xpack.security.enabled=false # 关闭安全认证,免密码访问
- xpack.security.http.ssl.enabled=false # 关闭 HTTPS 加密
- xpack.security.transport.ssl.enabled=false # 关闭节点传输加密
- ES_JAVA_OPTS=-Xms512m -Xmx512m # JVM 内存配置,避免占用过高
volumes:
- ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/es/data:/usr/share/elasticsearch/data
restart: always
# Kibana 最新稳定版:8.17.0(必须与 ES 版本完全一致)
kibana:
image: kibana:8.17.0
container_name: knowledge_hub_kibana
ports:
- "5601:5601" # Kibana 网页控制台端口
environment:
- ELASTICSEARCH_HOSTS=http://es:9200 # 连接 ES 容器内部地址
volumes:
- ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/kibana:/usr/share/kibana/data
restart: always
depends_on:
- es # 等待 ES 启动完成后再启动 Kibana
这里加一下 rabbitmq、es、kibana
(直接从代码仓库复制就行)
这里 es 用单独的 Dockerfile 跑,因为要加一下 IK 分词插件
# 官方 ES 基础镜像
FROM elasticsearch:8.17.0
# 安装 IK 分词(版本严格和 ES 一致)
RUN elasticsearch-plugin install --batch \
https://release.infinilabs.com/analysis-ik/stable/elasticsearch-analysis-ik-8.17.0.zip
跑一下:
Video: wxv_4638965120990412806
然后改下代码:
首先改一下 document 模块,加一个 publish 接口
然后加一个 mq 模块,用来封装 RabbitMQ 的发布订阅逻辑
也就是:
发布 → 投递 → 消费
之后加一个 pipeline 模块,用来监听 RabbitMQ 队列,收到消息后的流水线处理
也就是:
消费 → 分块 → 向量化 → 写入 kh_chunk

涉及到这三个模块:

我们具体看一下代码:
Video: wxv_4638967595428413441
这样我们就把 publish 接口到 RabbitMQ 到消费者处理消息
消费者的 分块 → 向量化 → 存到 ES kh_chunk 表的流程理清了
然后详细看一下 Chunk、Embedding、VectorIndex 这三块逻辑:
首先是 Chunk:
Video: wxv_4638965925625724930
chunk 拆分就是按照 markdown 格式来的,标题、段落等,拼接到目标块大小,直接用 langchain 的 RecursiveCharacterTextSplitter 来做
然后是 Embedding:
Video: wxv_4638973336289755140
这个也是直接调用 LangChain 的 OpenAIEmbeddings
这两步就可以体验用 Agent 框架的好处了,干啥都有对应的 API。
最后一步就是把向量存到 ES 的表里:
ElasticSearch 也支持向量字段:

指定下对应类型就可以。
这样我们直接用 ES 就能做向量+ 关键词混合检索。
Video: wxv_4638968932958355465
至此,整个流程就讲了一遍了:

然后我们跑一下:
根据这个 curl3.md 来:
---
## 发布文档(直接发布,无审核)
DOC_ID='docid'
将草稿(或已发布文档)设为已发布,并投递 RabbitMQ 异步管线:RAG 向量化。
curl -s -X PUT "http://localhost:3000/documents/${DOC_ID}/publish"
成功后 `status=1`,异步消费会执行:
- **RAG**:Markdown 分块 → Embedding → 写入 Elasticsearch `kh_chunk`(dense_vector)
---
GET /_cat/indices?
GET /kh_chunk/_search
{
"size": 100,
"query": {
"match_all": {}
}
}
记得在 .env 配一下环境变量:
OPENAI_API_KEY=sk-xxx
EMBEDDING_BASE_URL=https://dashscope.aliyuncs.com/compatible-mode/v1
EMBEDDING_MODEL=text-embedding-v3
试一下:
Video: wxv_4638969783629643778
测试没问题。
至此,从 publish 接口到 RabbitMQ 的消息队列。
然后消费者的 RAG 流水线:分块、向量化、存到 ES 索引表。
整个流程就都跑通了。
总结
这节我们实现了 publish 接口。
发布会修改文档状态,然后发一条消息到 RabbitMQ 消息队列。
消费者监听队列的消息,收到后会拿到文档,做文档分块、向量化、存到 ES 索引表的流水线
当然,这个过程不只是要做向量化,还要存到全文检索库、提取实体存到知识图谱,还有其他的消费者,下节继续实现。