You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 18:45:00