使用Haystack+Milvus搭建RAG Pipeline时的组件报错求助
Haystack 2.x + Milvus RAG Pipeline 错误修复方案
核心问题分析
在Python 3.10.12集群环境搭建RAG Pipeline时,遇到两类错误:
- 使用
@component装饰函数时,提示must have a 'run()' method:Haystack 2.x的@component仅支持装饰类,且类必须包含run()作为执行入口,直接装饰函数不符合规范。 - 移除装饰器后提示
doesn't seem to be a component:普通函数无法被识别为Haystack Component实例,不能加入Pipeline。
此外原代码还存在模型重复加载、Pipeline连接无效等问题。
解决步骤
1. 清理环境包冲突
当前环境同时存在Haystack 1.x(farm-haystack)和2.x(haystack-ai)包,会导致兼容性问题,需卸载旧版本:
pip uninstall -y farm-haystack haystack
仅保留haystack-ai和milvus-haystack即可。
2. 修正自定义组件实现
自定义组件必须以类的形式编写,通过@component装饰,类中包含:
__init__方法:初始化模型和Tokenizer,避免每次调用重复加载run方法:定义组件输入输出逻辑,返回{输出名: 输出值}格式的字典
自定义嵌入与生成组件代码
from haystack import component from transformers import AutoModelForSeq2SeqLM, AutoTokenizer import torch # 自定义文档嵌入组件 @component class ModelEmbedder: def __init__(self, model_name: str, cache_dir: str = None): self.tokenizer = AutoTokenizer.from_pretrained(model_name, cache_dir=cache_dir) self.model = AutoModelForSeq2SeqLM.from_pretrained(model_name, cache_dir=cache_dir) self.device = torch.device("cuda" if torch.cuda.is_available() else "cpu") self.model.to(self.device) def run(self, documents): embeddings = [] for doc in documents: inputs = self.tokenizer( doc.content, padding="max_length", truncation=True, return_tensors="pt" ).to(self.device) with torch.no_grad(): output = self.model(**inputs) # 若模型无pooler_output,替换为last_hidden_state均值: # embedding = output.last_hidden_state.mean(dim=1).squeeze(0).cpu().numpy() embedding = output.pooler_output.squeeze(0).cpu().numpy() embeddings.append(embedding) for doc, embedding in zip(documents, embeddings): doc.embedding = embedding return {"documents": documents} # 自定义回答生成组件 @component class ModelGenerator: def __init__(self, model_name: str, cache_dir: str = None): self.tokenizer = AutoTokenizer.from_pretrained(model_name, cache_dir=cache_dir) self.model = AutoModelForSeq2SeqLM.from_pretrained(model_name, cache_dir=cache_dir) self.device = torch.device("cuda" if torch.cuda.is_available() else "cpu") self.model.to(self.device) def run(self, query: str, documents: list = None, generation_kwargs: dict = {}): # 构建带上下文的输入文本 context_text = "\n".join([doc.content for doc in documents]) if documents else "" input_text = f"问题:{query}\n上下文:{context_text}\n回答:" inputs = self.tokenizer( input_text, padding="max_length", truncation=True, return_tensors="pt" ).to(self.device) with torch.no_grad(): output = self.model.generate(**inputs, **generation_kwargs) generated_text = self.tokenizer.decode(output[0], skip_special_tokens=True) return {"response": generated_text}
3. 修复Pipeline连接逻辑
原Pipeline存在无效连接(如未定义的text_embedder),需补充查询嵌入组件并修正数据流向:
修正后的RAG Pipeline代码
from haystack import Pipeline from haystack.components.converters import MarkdownToDocument from haystack.components.preprocessors import DocumentSplitter from haystack.components.writers import DocumentWriter from haystack.components.builders import PromptBuilder from milvus_haystack import MilvusDocumentStore from milvus_haystack.milvus_embedding_retriever import MilvusEmbeddingRetriever # 初始化Milvus文档存储(根据实际集群配置修改参数) document_store = MilvusDocumentStore( host="your-milvus-host", port="19530", collection_name="rag_collection", embedding_dim=768 # 替换为你的模型嵌入维度 ) # 定义提示模板 prompt_template = """ 根据以下上下文回答问题: {% for doc in documents %} {{ doc.content }} {% endfor %} 问题:{{ query }} 回答: """ # 构建Pipeline rag_pipeline = Pipeline() # 数据导入流程组件 rag_pipeline.add_component("converter", MarkdownToDocument()) rag_pipeline.add_component("splitter", DocumentSplitter(split_by="sentence", split_length=2)) rag_pipeline.add_component("embedder", ModelEmbedder(model_name=mymodel, cache_dir=cache_dir)) rag_pipeline.add_component("writer", DocumentWriter(document_store=document_store)) # 查询流程组件 rag_pipeline.add_component("query_embedder", ModelEmbedder(model_name=mymodel, cache_dir=cache_dir)) rag_pipeline.add_component("retriever", MilvusEmbeddingRetriever(document_store=document_store, top_k=3)) rag_pipeline.add_component("prompt_builder", PromptBuilder(template=prompt_template)) rag_pipeline.add_component("generator", ModelGenerator(model_name=mymodel, cache_dir=cache_dir)) # 连接数据导入流程 rag_pipeline.connect("converter.documents", "splitter.documents") rag_pipeline.connect("splitter.documents", "embedder.documents") rag_pipeline.connect("embedder.documents", "writer.documents") # 连接查询流程 rag_pipeline.connect("query_embedder.documents", "retriever.query_embedding") rag_pipeline.connect("retriever.documents", "prompt_builder.documents") rag_pipeline.connect("prompt_builder.prompt", "generator.query") # 绘制Pipeline结构 rag_pipeline.draw('./rag_pipeline.png')
注意事项
- 替换
mymodel为实际的Hugging Face模型名,若模型无pooler_output,需调整嵌入获取逻辑(如使用last_hidden_state均值) - 根据Milvus集群实际配置修改
MilvusDocumentStore的参数 - 集群环境有GPU时,确保
torch能正确识别以加速模型运行
内容的提问来源于stack exchange,提问作者ArieAI
相关产品推荐
相关产品推荐

