基于 MLflow 与 LlamaIndex Workflow 构建混合检索 RAG 应用:从构建、记录到评估与追踪的完整实战
【免费下载链接】mlflowThe open source AI engineering platform for agents, LLMs, and ML models. MLflow enables teams of all sizes to debug, evaluate, monitor, and optimize production-quality AI applications while controlling costs and managing access to models and data.项目地址: https://gitcode.com/GitHub_Trending/ml/mlflow
本文是一篇围绕仓库中 examples/llama_index/workflow 示例目录展开的实战指南,核心场景是:使用 LlamaIndex 的事件驱动 Workflow 编排框架搭建一个同时融合向量检索(Vector Search)、BM25 关键词检索与 Web 搜索的混合检索 RAG 问答系统,再借助 MLflow 的模型记录(Model-from-code)、Experiment 跟踪、Evaluate 评估与Tracing 追踪能力,实现从构建、实验对比到质量诊断的完整闭环。阅读本文后,你将掌握如何用代码定义可配置的 RAG Workflow、如何用一条命令将其记录为可复现的 MLflow 模型、如何对多种检索策略做量化评估,以及如何通过 Trace UI 定位回答质量问题的根因。
为什么要做混合检索:单一检索方式的局限
Retrieval-Augmented Generation(RAG)通过给 LLM 注入外部知识来提升回答质量,但检索环节常常成为瓶颈:基于 embedding 的向量检索未必总能命中语义最相关的内容(例如专有名词、缩写、精确术语),而纯关键词检索又缺乏语义理解。业界虽有大量检索增强技巧,却不存在放之四海皆准的单一方案。
因此,一个务实策略是把多种检索方法并行组合:向量检索负责语义召回,BM25 负责精确关键词命中,Web 搜索负责补充知识库之外的时效性信息;随后将多路结果合并、去重并重排,过滤掉无关内容后再交给 LLM 作答。本示例正是按这一思路设计的。
LlamaIndex Workflow:事件驱动的编排基础
在深入代码前,先理解 LlamaIndex Workflow 的三大核心抽象(来源:教程笔记 与 workflow.py):
- Step(步骤):执行单元,代表工作流中的一个具体动作,以
@step装饰的异步方法实现; - Event(事件):触发 Step 的信号,在 Step 之间传递数据、控制流转方向;
- Workflow(工作流):把二者连接成一个 Python 类,每个 Step 是类的成员方法,显式声明输入与输出事件。
这种"事件驱动"设计天然适合并行/异步执行:多个检索 Step 可以同时触发、互不阻塞,再通过上下文对象统一汇聚结果,这为处理耗时检索任务与生产级扩展性提供了基础。
环境准备与依赖安装
示例目录包含完整的 Workflow 定义(workflow/子目录)、手把手教学笔记 Tutorial.ipynb 以及样本数据集(mlflow_qa_dataset.csv、urls.txt)。
克隆仓库并安装依赖
按照 README.md 中的说明,克隆仓库后进入示例目录并运行安装脚本:
git clone https://github.com/mlflow/mlflow.gitcd mlflow/examples/llama_index/workflow chmod +x install.sh ./install.sh安装完成后,在 Poetry 环境内启动 Jupyter Notebook:
poetry run jupyter notebook安装脚本做了什么
install.sh 的核心逻辑如下:
- 将
$HOME/.local/bin加入 PATH;若未检测到 Poetry 则通过官方安装脚本自动安装; - 检查 Docker:由于向量库 Qdrant 通常以 Docker 方式运行,脚本会检测 Docker 是否可用并给出提示;若你不需要 Qdrant(例如只跑 BM25/Web 检索),可加
--no-qdrant标志跳过该检查; - 通过
poetry run pip install jupyter安装 Jupyter Notebook; - 执行
poetry install依据 pyproject.toml 安装全部依赖。
依赖清单解读
pyproject.toml 声明了本示例的完整依赖栈(Python 版本要求>=3.10,<3.13):
mlflow >= 2.17.0:提供实验跟踪、模型记录、评估与追踪能力;llama-index >= 0.11.0:Workflow 编排框架核心;llama-index-postprocessor-rankgpt-rerank:基于 RankGPT 的重排后处理器;llama-index-readers-web:SimpleWebPageReader,用于加载网页文档;llama-index-retrievers-bm25:BM25 关键词检索器;llama-index-tools-tavily-research:Tavily Web 搜索工具(LLM 场景优化的搜索 API);llama-index-utils-workflow:Workflow 工具集;llama-index-vector-stores-qdrant:Qdrant 向量存储集成。
起步:创建实验、配置 LLM 与 Embedding
创建 MLflow Experiment
Experiment是 MLflow 组织模型开发过程的最小单元,记录模型定义、配置、参数、依赖版本等信息。在笔记中创建新实验:
import mlflow mlflow.set_experiment("LlamaIndex Workflow RAG")配置 LLM 与 Embedding
LlamaIndex 通过全局Settings对象统一管理 LLM 与 Embedding,后续所有 LlamaIndex 组件都会使用这里的配置。MLflow 在记录模型时会自动把Settings配置写入实验,从而保证跨环境的可复现性。
方案一:OpenAI(默认)
LlamaIndex 默认使用 OpenAI API。示例代码(截至 2024 年 10 月默认模型为gpt-3.5-turbo与text-embeddings-ada-002)推荐换用更新、更高效的模型以获得更好效果与更低成本:
import getpass import os os.environ["OPENAI_API_KEY"] = getpass.getpass("Enter your OpenAI API key")from llama_index.core import Settings from llama_index.embeddings.openai import OpenAIEmbedding from llama_index.llms.openai import OpenAI Settings.embed_model = OpenAIEmbedding(model="text-embedding-3-small") Settings.llm = OpenAI(model="gpt-4o-mini")方案二:其他托管模型(以 Databricks 托管的 Llama3.1 70B 为例)
- 安装对应提供方的集成包(如
llama-index-llms-databricks); - 按集成文档设置所需环境变量(如
DATABRICKS_SERVING_ENDPOINT与DATABRICKS_TOKEN); - 实例化 LLM 与 Embedding 并写入
Settings:
from llama_index.core import Settings from llama_index.embeddings.databricks import DatabricksEmbedding from llama_index.llms.databricks import Databricks Settings.embed_model = DatabricksEmbedding(model="databricks-gte-large-en") Settings.llm = Databricks(model="databricks-meta-llama-3-1-70b-instruct")方案三:本地模型:LlamaIndex 同样支持本地部署的 LLM,按 LlamaIndex 本地模型入门教程配置即可。
配置 Web 搜索 API
本示例的 Web 检索使用Tavily AI——为 LLM 应用优化的搜索 API,与 LlamaIndex 原生集成。访问其官网申请免费额度后设置环境变量(也可换成 LlamaIndex 支持的其他搜索引擎,如 Google Search Tool):
import getpass import os os.environ["TAVILY_AI_API_KEY"] = getpass.getpass("Enter your Tavily AI APi Key")构建检索索引:向量索引与 BM25 索引
加载文档
urls.txt 中存放了一批 MLflow 官方文档页面地址,通过SimpleWebPageReader加载为 LlamaIndex 文档对象:
from llama_index.readers.web import SimpleWebPageReader with open("data/urls.txt") as file: urls = [line.strip() for line in file if line.strip()] documents = SimpleWebPageReader(html_to_text=True).load_data(urls)向量索引(Qdrant)
将文档写入向量数据库。教程选用Qdrant(自托管免费),先用 Docker 启动服务:
$ docker pull qdrant/qdrant $ docker run -p 6333:6333 -p 6334:6334 \ -v $(pwd)/.qdrant_storage:/qdrant/storage:z \ qdrant/qdrant然后创建连接 Qdrant 的索引对象并摄入文档:
import qdrant_client from llama_index.vector_stores.qdrant import QdrantVectorStore client = qdrant_client.QdrantClient(host="localhost", port=6333) vector_store = QdrantVectorStore(client=client, collection_name="mlflow_doc") from llama_index.core import StorageContext, VectorStoreIndex storage_context = StorageContext.from_defaults(vector_store=vector_store) index = VectorStoreIndex.from_documents(documents=documents, storage_context=storage_context)当然也可以换成 FAISS、Chroma、Databricks Vector Search 等 LlamaIndex 支持的任意向量库——若更换,需要按对应文档调整 workflow.py 中的向量存储接入代码。
BM25 关键词索引
BM25 检索基于本地节点文件,先对文档分块(chunk_size=512)后构建检索器并持久化到.bm25_retriever目录,供工作流运行时加载:
from llama_index.core.node_parser import SentenceSplitter from llama_index.retrievers.bm25 import BM25Retriever splitter = SentenceSplitter(chunk_size=512) nodes = splitter.get_nodes_from_documents(documents) bm25_retriever = BM25Retriever.from_defaults(nodes=nodes) bm25_retriever.persist(".bm25_retriever")深入工作流实现:events、prompts 与 workflow 类
workflow/目录包含三份核心代码:events.py(事件定义)、prompts.py(提示词模板)与 workflow.py(主工作流类)。
事件定义:数据如何在步骤间流转
events.py 中的事件都是携带字段的 Pydantic 模型。例如VectorSearchRetrieveEvent携带用户查询,触发向量检索步骤:
class VectorSearchRetrieveEvent(Event): """Event for triggering VectorStore index retrieval step.""" query: str完整的 7 种事件及其职责:
| 事件 | 携带字段 | 职责 |
|---|---|---|
VectorSearchRetrieveEvent | query | 触发向量库索引检索步骤 |
BM25RetrieveEvent | query | 触发 BM25 检索步骤 |
TransformQueryEvent | query | 将用户问题改写为搜索友好查询 |
WebsearchEvent | search_query | 触发 Web 搜索工具步骤 |
RetrievalResultEvent | nodes,retriever | 把各路检索结果(含来源标识)送回汇聚步骤 |
RerankEvent | nodes | 把合并后的节点送往重排步骤 |
QueryEvent | context | 触发最终问答步骤 |
注意RetrievalResultEvent.retriever的类型被限定为Literal["vector_search", "bm25", "web_search"],与工作流支持的三类检索器一一对应。
提示词模板
prompts.py 定义了两处 LLM 调用提示:
TRANSFORM_QUERY_TEMPLATE:把用户原始问题提炼成更适合搜索引擎的查询串("只输出优化后的查询");FINAL_QUERY_TEMPLATE:要求 LLM只依据给定的上下文作答、不得使用先验知识("只输出答案")。
工作流类:可配置的混合检索编排
workflow.py 定义了HybridRAGWorkflow(Workflow)。构造函数通过retrievers参数声明要启用的检索方式,支持{"vector_search", "bm25", "web_search"}的任意子集:
class HybridRAGWorkflow(Workflow): VALID_RETRIEVERS = {"vector_search", "bm25", "web_search"} def __init__(self, retrievers=None, **kwargs): super().__init__(**kwargs) self.llm = Settings.llm self.retrievers = retrievers or [] if invalid_retrievers := set(self.retrievers) - self.VALID_RETRIEVERS: raise ValueError(f"Invalid retrievers specified: {invalid_retrievers}") self._use_vs_retriever = "vector_search" in self.retrievers self._use_bm25_retriever = "bm25" in self.retrievers self._use_web_search = "web_search" in self.retrievers if self._use_vs_retriever: qd_client = qdrant_client.QdrantClient(host=_QDRANT_HOST, port=_QDRANT_PORT) vector_store = QdrantVectorStore(client=qd_client, collection_name=_QDRANT_COLLECTION_NAME) index = VectorStoreIndex.from_vector_store(vector_store=vector_store) self.vs_retriever = index.as_retriever() if self._use_bm25_retriever: self.bm25_retriever = BM25Retriever.from_persist_dir(_BM25_PERSIST_DIR) if self._use_web_search: self.tavily_tool = TavilyToolSpec(api_key=os.environ.get("TAVILY_AI_API_KEY"))文件头部定义了外部依赖的环境变量,均有默认值,便于本地快速启动:
QDRANT_HOST(默认localhost)QDRANT_PORT(默认6333)QDRANT_COLLECTION_NAME(默认mlflow_doc)
动态决定检索器是本设计的精髓:只需改变retrievers列表,就能实验不同检索组合,而不必为几乎相同的模型代码做多份复制。
第一步:路由检索(route_retrieval)
工作流以接收StartEvent的route_retrieval步骤为入口,它把查询写入上下文,然后根据配置并行派发事件——这是事件驱动框架实现并行异步执行的关键:
# If no retriever is specified, proceed directly to the final query step with an empty context if len(self.retrievers) == 0: return QueryEvent(context="") # Trigger the retrieval steps based on the configuration if self._use_vs_retriever: ctx.send_event(VectorSearchRetrieveEvent(query=query)) if self._use_bm25_retriever: ctx.send_event(BM25RetrieveEvent(query=query)) if self._use_web_search: ctx.send_event(TransformQueryEvent(query=query))当retrievers为空时,工作流退化为"仅凭 LLM 先验知识直接作答",这正好可以作为混合检索的对照基线。
三类检索步骤
向量检索与 BM25 检索都很直接——调用各自的retrieve()并包装成RetrievalResultEvent:
@step async def query_vector_store(self, ev: VectorSearchRetrieveEvent) -> RetrievalResultEvent: nodes = self.vs_retriever.retrieve(ev.query) return RetrievalResultEvent(nodes=nodes, retriever="vector_search") @step async def query_bm25(self, ev: BM25RetrieveEvent) -> RetrievalResultEvent: nodes = self.bm25_retriever.retrieve(ev.query) return RetrievalResultEvent(nodes=nodes, retriever="bm25")Web 检索则多一步:先用 LLM 把原始问题改写为适合搜索引擎的查询串(transform_query),再调用 Tavily 工具执行搜索(query_web_search,max_results=5):
@step async def transform_query(self, ev: TransformQueryEvent) -> WebsearchEvent: prompt = TRANSFORM_QUERY_TEMPLATE.format(query=ev.query) transformed_query = self.llm.complete(prompt).text return WebsearchEvent(search_query=transformed_query)汇聚与重排(gather_retrieval_results/rerank)
gather_retrieval_results用ctx.collect_events()异步轮询各检索步骤的结果——在工作流框架下,只要还有检索器未返回,collect_events就返回None,步骤等待下一轮轮询;全部到齐后再决定走哪条路:
results = ctx.collect_events(ev, [RetrievalResultEvent] * len(self.retrievers))- 只有一个检索器:跳过重排,直接把节点文本拼接为上下文,进入最终问答;
- 多个检索器:合并所有节点,并给每个节点打上来源标记(
node.metadata["retriever"]),交给rerank步骤。
多路结果直接拼接会造成上下文过大、且混入无关/重复内容。由于 Web 搜索结果没有相似度分数,常见的分数排序方案行不通,因此这里改用LLM 重排:借助 RankGPT 集成的RankGPTRerank按查询相关性排序并取前 5:
reranker = RankGPTRerank(llm=self.llm, top_n=5) reranked_nodes = reranker.postprocess_nodes(ev.nodes, query_str=query) reranked_context = "\n".join(node.text for node in reranked_nodes)最终问答(query_result)
把重排后的上下文与用户查询填入FINAL_QUERY_TEMPLATE,调用 LLM 生成答案,以StopEvent(result=...)结束整个工作流:
@step async def query_result(self, ctx: Context, ev: QueryEvent) -> StopEvent: query = await ctx.get("query") prompt = FINAL_QUERY_TEMPLATE.format(context=ev.context, query=query) response = self.llm.complete(prompt).text return StopEvent(result=response)实例化并运行
以"向量检索 + BM25"组合为例(timeout=60秒),在 Jupyter 中直接异步运行:
from workflow.workflow import HybridRAGWorkflow workflow = HybridRAGWorkflow(retrievers=["vector_search", "bm25"], timeout=60) response = await workflow.run(query="Why use MLflow with LlamaIndex?") print(response)用 Model-from-code 把工作流记录进 MLflow
运行评估前,先把工作流作为模型记录到 MLflow 中。本示例采用Model-from-code方式:模型以独立 Python 脚本的形式记录,代码是模型定义的唯一事实来源,规避了 pickle 等序列化方式的稳定性风险,再结合 MLflow 的依赖环境冻结能力,实现可靠的模型持久化。
模型入口脚本
workflow/model.py 是记录模型的入口:它通过mlflow.models.ModelConfig()单例读取记录时传入的model_config,从而用同一份代码实例化不同配置的工作流,无需重复模型定义:
from workflow.workflow import HybridRAGWorkflow import mlflow # Get model config from ModelConfig singleton model_config = mlflow.models.ModelConfig() retrievers = model_config.get("retrievers") # Create the workflow instance. workflow = HybridRAGWorkflow(retrievers=retrievers, timeout=300) # Set the model instance logging. This is mandatory for using model-from-code logging method. mlflow.models.set_model(workflow)其中mlflow.models.set_model(workflow)是 Model-from-code 方式记录模型的必选项。
用 mlflow.llama_index.log_model 记录多套配置
mlflow.llama_index.log_model 支持记录 Index、Engine、Workflow 对象,或指向包含上述类型定义脚本的路径(字符串);配合model_config参数即可为同一脚本注入不同配置。此处为 4 种检索策略各开一个 Run 并记录模型:
# 1. No retrievers (prior knowledge in LLM). # 2. Vector search retrieval only. # 3. Vector search and keyword search (BM25) # 4. All retrieval methods including web search. run_name_to_retrievers = { "none": [], "vs": ["vector_search"], "vs + bm25": ["vector_search", "bm25"], "vs + bm25 + web": ["vector_search", "bm25", "web_search"], } models = [] for run_name, retrievers in run_name_to_retrievers.items(): with mlflow.start_run(run_name=run_name): model_info = mlflow.llama_index.log_model( # Specify the model Python script. llama_index_model="workflow/model.py", # Specify retrievers to use. model_config={"retrievers": retrievers}, # Define dependency files to save along with the model code_paths=["workflow"], # Subdirectory to save artifacts (not important) name="model", ) models.append(model_info)关键参数说明(以 model.py 的签名与文档为准):
llama_index_model:模型对象或模型脚本路径;model_config:记录时保存的配置,Model-from-code 方式下通过ModelConfig对象在模型代码内读取,实现配置与代码解耦;code_paths:随模型一起保存的依赖代码目录,本示例传入["workflow"]以确保workflow包可被加载;name:artifacts 保存的子目录名。
记录完成后打开 MLflow UI,可以看到 4 个 Run 以不同retrievers参数值被记录下来;点击 Run 名并进入 "Artifacts" 页签,可查看模型文件、依赖版本与Settings配置等元数据。值得留意的是,MLflow 记录Settings时会刻意跳过 API Key(避免密钥泄漏)与不可序列化的函数对象。
用 mlflow.evaluate 量化评估检索策略
启用 LlamaIndex Tracing
在评估前先开启MLflow Tracing,一行命令即可让 MLflow 自动追踪每一次 LlamaIndex 执行(其效果下一节详述):
mlflow.llama_index.autolog()从源码看,mlflow/llama_index/autolog.py 目前仅支持 tracing 自动记录:log_traces=True时安装 LlamaIndex tracer,disable/silent分别用于禁用与静默模式。
加载评估数据集
示例仓库自带含30 组问答对的评估数据 mlflow_qa_dataset.csv,字段为query与ground_truth,问题覆盖 MLflow 的各类用法(例如"如何用 Amazon S3 存储 MLflow 模型 artifacts"及其标准答案):
import pandas as pd eval_df = pd.read_csv("data/mlflow_qa_dataset.csv") display(eval_df.head(3))执行评估
mlflow.evaluate()需要三要素:数据集、已记录的模型、要计算的指标。本示例对每种配置分别评估,采用两项指标:
- Latency(
latency()):单条查询执行整个工作流所耗时间; - Answer Correctness(
answer_correctness("openai:/gpt-4o-mini")):基于ground_truth,由 OpenAI GPT-4o 模型按 1–5 分打分,衡量答案正确性。
from mlflow.metrics import latency from mlflow.metrics.genai import answer_correctness for model_info in models: with mlflow.start_run(run_id=model_info.run_id): result = mlflow.evaluate( # Pass the URI of the logged model above model=model_info.model_uri, data=eval_df, # Specify the column for ground truth answers. targets="ground_truth", # Define the metrics to compute. extra_metrics=[ latency(), answer_correctness("openai:/gpt-4o-mini"), ], # The answer_correctness metric requires "inputs" column to be # present in the dataset. We have "query" instead so need to # specify the mapping in `evaluator_config` parameter. evaluator_config={"col_mapping": {"inputs": "query"}}, )注意answer_correctness需要数据集中存在inputs列,而本数据集列名是query,因此通过evaluator_config={"col_mapping": {"inputs": "query"}}完成列映射。这两项指标仅为演示,你完全可以补充 toxicity、faithfulness 等指标或自定义指标。
阅读评估结果
评估耗时几分钟,完成后在 MLflow UI 的 Experiment 页面点击 Run 列表上方的图表图标,即可看到对比结果:
第一行是 answer correctness 柱状图,第二行是 latency 结果。本例中最优组合是"Vector Search + BM25";有趣的是,加入 Web 搜索后不仅延迟显著上升,answer correctness 反而下降——例如在回答"如何启动 Model Registry"时,开启 Web 搜索的模型给出了关于模型部署的离题答案,而vs + bm25组合回答正确。(评估结果会因模型配置与随机性略有差异。)
用 MLflow Tracing 定位回答质量问题的根因
由于唯一变化的是检索策略,问题大概率出在检索环节,但仅看最终答案很难判断各路检索器各自返回了什么。此时MLflow Tracing派上用场:它完整记录工作流执行期间每一步的输入、输出、元数据与延迟,并与 LlamaIndex 深度集成。
在 Experiment 页面切换到"Traces" 页签,找到请求列为该问题、Run 名为 "vs + bm25 + web" 的记录,点击请求 ID 即可打开 Trace UI 查看各步骤详情:
本案例中,通过检查rerank(重排)步骤就锁定了问题:Web 检索器返回了与模型服务相关的无关上下文,而重排器错误地把它排为最相关。有了这一洞察,就可以针对性地改进——例如优化重排器对 MLflow 主题的理解、提高 Web 检索精度,甚至直接移除 Web 检索器。
小结与延伸
本示例完整展示了 LlamaIndex 与 MLflow 组合如何提升 RAG 工作流的开发效率与可观测性:
- Experiment Tracking:以 Run 维度组织并记录不同工作流配置,保证可复现性,支持跨 Run 性能追踪;
- MLflow Evaluate:对多种检索策略(无检索、仅向量、向量+BM25、全量混合)用 latency 与 answer correctness 做量化对比;
- MLflow UI:直观可视化不同检索策略对准确率与延迟的影响,帮助选出最优配置;
- MLflow Tracing:与 LlamaIndex 集成,提供工作流每一步的细粒度可观测性,用于诊断诸如重排错误之类的质量问题。
由此,你获得了一整套"构建 → 记录 → 评估 → 诊断 → 优化"的 RAG 开发闭环。进一步探索的仓库入口:
- 完整教学笔记:Tutorial.ipynb;
- 工作流核心实现:workflow.py、events.py、prompts.py、model.py;
- 环境与依赖:install.sh、pyproject.toml;
- 数据与示例图:mlflow_qa_dataset.csv、urls.txt;
- MLflow 集成源码:mlflow/llama_index/model.py、mlflow/llama_index/autolog.py。
【免费下载链接】mlflowThe open source AI engineering platform for agents, LLMs, and ML models. MLflow enables teams of all sizes to debug, evaluate, monitor, and optimize production-quality AI applications while controlling costs and managing access to models and data.项目地址: https://gitcode.com/GitHub_Trending/ml/mlflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考