实用指南:开源RAG工具终极技术手册
《实用笔记操作指南:开源RAG工具终极技术手册》——专为采用该架构的团队提供的合同规范、校验机制及可直接插入的代码模块。
本指南将逐步展示如何从原始材料构建出可使用的系统,内容来自《开源RAG工具终极技术指南:Docling、LlamaIndex、LangChain、Haystack与RAGAS》。重点在于具体的操作步骤、明确的检查点,以及可直接放入代码库的代码,无需猜测其用途。
2026年构建可用于生产的检索增强生成系统
在构建可投入生产的检索增强生成阶段,应在修改代码之前明确输入内容、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。除了功能结果外,还需记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。必须注明实际用于生成答案的对应内容。如果没有引用依据,操作人员就无法区分幻觉内容与索引缺失的问题。
前提条件
在准备阶段,应在修改代码之前明确输入内容、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能标志应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审核。 需引用实际作为答案依据的段落。如果没有引用,操作人员就无法区分是虚假信息还是索引缺失导致的错误。
1. Docling:智能文档处理基础平台
在1 Docling智能文档阶段,修改代码之前需明确输入内容、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都是产品本身的组成部分,而非后续需要补充的功能。 需引用实际作为答案依据的段落。如果没有引用,操作人员就无法区分是幻觉内容还是索引缺失导致的错误。
安装与基本使用
在安装与基础使用阶段,应在修改代码之前明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。相比冗长的脚本,更应采用小型且可测试的单元。当某个步骤失败时,故障原因应指向单一责任模块,而非复杂的流程链。需引用实际作为答案依据的段落;没有引用的话,操作人员就无法区分是虚假信息还是索引缺失所致。
pip install docling
from docling.document_converter import DocumentConverter
# Convert a PDF document to structured Markdown
converter = DocumentConverter()
result = converter.convert("sample_document.pdf")
structured_markdown = result.document.export_to_markdown()
print(structured_markdown)
利用OCR进行高级文档处理
对于带有阶段的高级文档处理,应在修改代码之前明确输入内容、各步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 应将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,定义成功判定标准,并拒绝默许的半完成状态。 需引用实际作为答案依据的段落。没有引用的话,操作人员就无法区分幻觉内容与索引缺失问题。
from docling.document_converter import DocumentConverter
from docling.datamodel.pipeline_options import PipelineOptions
# Configure pipeline with OCR for scanned documents
pipeline_options = PipelineOptions(
do_ocr=True, # Enable OCR for scanned documents
ocr_engine="tesseract", # Use Tesseract OCR engine
ocr_language="eng" # English language
)
converter = DocumentConverter(pipeline_options=pipeline_options)
result = converter.convert("scanned_document.pdf")
# Extract structured data including tables and images
document = result.document
tables = document.tables
images = document.images
print(f"Found {len(tables)} tables and {len(images)} images")
与其他框架的集成
在“与其他框架集成”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及完成标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。除了功能结果外,还需记录执行时间以及令牌或查询成本。提前了解成本情况可以避免在从演示环境过渡到共享环境时出现意外费用。必须注明支撑答案的具体内容出处;没有引用的话,操作人员就无法区分是幻觉内容还是索引缺失导致的错误。
# Example: Integrating Docling with LangChain
from langchain_community.document_loaders import TextLoader
from docling.document_converter import DocumentConverter
def docling_to_langchain_docs(file_path):
"""Convert document using Docling and return as LangChain documents"""
converter = DocumentConverter()
result = converter.convert(file_path)
content = result.document.export_to_markdown()
# Create LangChain document
loader = TextLoader(file_path)
docs = loader.load()
docs[0].page_content = content
docs[0].metadata["source"] = file_path
return docs
# Usage
langchain_docs = docling_to_langchain_docs("technical_manual.pdf")
2. LlamaIndex:检索优先架构
在2层LlamaIndex检索优先架构阶段,修改代码之前需明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储和功能标志应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审核。 需注明实际作为答案依据的段落。若没有引用,操作人员就无法区分是幻觉内容还是索引缺失导致的错误。
基础RAG实现
在基础RAG实现阶段,应在修改代码之前明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都是产品的一部分,而非后续需要完善的内容。 必须注明实际用于支撑答案的原文段落。如果没有引用依据,操作人员就无法区分幻觉内容与索引缺失问题。
from llama_index.core import VectorStoreIndex, SimpleDirectoryReader
from llama_index.llms.openai import OpenAI
# Load documents
documents = SimpleDirectoryReader("./data").load_data()
# Create index
index = VectorStoreIndex.from_documents(documents)
# Create query engine
query_engine = index.as_query_engine(
similarity_top_k=3,
response_mode="compact"
)
# Query the system
response = query_engine.query("What are the main features of the product?")
print(response)
带自动合并的层次化分块
在“分层分块与自动合并”阶段,修改代码之前需明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 优先选择小型、可测试的单元,而非冗长的脚本。当某个步骤失败时,故障应指向单一责任点,而非复杂的流程链。 需引用实际作为答案依据的段落。没有引用的话,操作人员就无法区分是幻觉内容还是索引缺失所致。
from llama_index.core import VectorStoreIndex
from llama_index.core.node_parser import HierarchicalNodeParser, get_leaf_nodes
from llama_index.core.retrievers import AutoMergingRetriever
from llama_index.core.query_engine import RetrieverQueryEngine
# Create hierarchical nodes: 2048 -> 512 -> 128 token chunks
node_parser = HierarchicalNodeParser.from_defaults(
chunk_sizes=[2048, 512, 128]
)
nodes = node_parser.get_nodes_from_documents(documents)
leaf_nodes = get_leaf_nodes(nodes)
# Build index on leaf nodes only
index = VectorStoreIndex(leaf_nodes)
index.storage_context.docstore.add_documents(nodes)
# Auto-merging retriever replaces small chunks with parent context when relevant
retriever = AutoMergingRetriever(
index.as_retriever(similarity_top_k=3),
index.storage_context,
merge_batch_size=5,
)
query_engine = RetrieverQueryEngine(retriever)
response = query_engine.query("Explain the technical specifications in detail")
print(response)
使用 LlamaIndex 进行自定义评估
在基于LlamaIndex的定制评估阶段,修改代码之前需明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 应将此阶段视为输入与验证后输出之间的契约。为相关成果命名,设定成功检测标准,并拒绝默许的半完成状态。 需注明实际作为答案依据的段落。若没有引用,操作人员就无法区分幻觉内容与索引缺失问题。
from llama_index.core.evaluation import FaithfulnessEvaluator, RelevancyEvaluator
from llama_index.core import Settings
# Initialize evaluators
faithfulness_evaluator = FaithfulnessEvaluator(llm=Settings.llm)
relevancy_evaluator = RelevancyEvaluator(llm=Settings.llm)
# Evaluate response
eval_result = faithfulness_evaluator.evaluate_response(
query="What are the system requirements?",
response=response,
contexts=[node.text for node in response.source_nodes]
)
print(f"Faithfulness score: {eval_result.score}")
print(f"Relevancy score: {relevancy_evaluator.evaluate_response(query='What are the system requirements?', response=response, contexts=[node.text for node in response.source_nodes]).score}")
3. LangChain:以工作流为中心的框架
在基于LangChain的以工作流为中心的框架阶段,修改代码之前需明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。
基础RAG流程
from langchain_community.document_loaders import PyPDFLoader
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_community.vectorstores import Chroma
from langchain_openai import OpenAIEmbeddings, ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnablePassthrough
from langchain_core.output_parsers import StrOutputParser
# Load and split documents
loader = PyPDFLoader("sample.pdf")
docs = loader.load()
text_splitter = RecursiveCharacterTextSplitter(chunk_size=1000, chunk_overlap=200)
splits = text_splitter.split_documents(docs)
# Create vector store
vectorstore = Chroma.from_documents(documents=splits, embedding=OpenAIEmbeddings())
# Create retriever
retriever = vectorstore.as_retriever()
# Create prompt template
template = """Answer the question based only on the following context:
{context}
Question: {question}
"""
prompt = ChatPromptTemplate.from_template(template)
# Create chain
llm = ChatOpenAI(model_name="gpt-4o", temperature=0)
rag_chain = (
{"context": retriever, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
# Execute
response = rag_chain.invoke("What are the key benefits mentioned in the document?")
print(response)
具有记忆功能的多步骤工作流
from langchain_core.messages import HumanMessage, AIMessage
from langchain_core.chat_history import BaseChatMessageHistory
from langchain_core.runnables.history import RunnableWithMessageHistory
from langchain_community.chat_message_histories import ChatMessageHistory
# Create chat history
class InMemoryHistory(BaseChatMessageHistory):
def __init__(self):
self.messages = []
def add_user_message(self, message: str):
self.messages.append(HumanMessage(content=message))
def add_ai_message(self, message: str):
self.messages.append(AIMessage(content=message))
def clear(self):
self.messages = []
# Store chat history
chat_histories = {}
def get_chat_history(session_id: str) -> BaseChatMessageHistory:
if session_id not in chat_histories:
chat_histories[session_id] = InMemoryHistory()
return chat_histories[session_id]
# Create chain with memory
chain_with_memory = RunnableWithMessageHistory(
rag_chain,
get_chat_history,
input_messages_key="question",
history_messages_key="chat_history",
)
# Execute with memory
session_id = "user_123"
response = chain_with_memory.invoke(
"What are the key benefits mentioned in the document?",
config={"configurable": {"session_id": session_id}}
)
print(response)
# Follow-up question
follow_up_response = chain_with_memory.invoke(
"Can you elaborate on the second benefit?",
config={"configurable": {"session_id": session_id}}
)
print(follow_up_response)
4. Haystack:生产级搜索引擎
基本管道配置
from haystack import Pipeline
from haystack.components.embedders import SentenceTransformersTextEmbedder
from haystack.components.retrievers import InMemoryBM25Retriever, InMemoryEmbeddingRetriever
from haystack.components.joiners import JoinDocuments
from haystack.components.generators import OpenAIGenerator
from haystack.components.preprocessors import DocumentCleaner, DocumentSplitter
from haystack.document_stores.in_memory import InMemoryDocumentStore
from haystack.utils import Secret
# Initialize components
document_store = InMemoryDocumentStore()
cleaner = DocumentCleaner()
splitter = DocumentSplitter(split_by="word", split_length=1000)
embedder = SentenceTransformersTextEmbedder(model="sentence-transformers/all-MiniLM-L6-v2")
bm25_retriever = InMemoryBM25Retriever(document_store=document_store)
embedding_retriever = InMemoryEmbeddingRetriever(document_store=document_store)
joiner = JoinDocuments(join_mode="concatenate")
generator = OpenAIGenerator(api_key=Secret.from_env_var("OPENAI_API_KEY"), model="gpt-4o")
# Create pipeline
pipeline = Pipeline()
pipeline.add_component("cleaner", cleaner)
pipeline.add_component("splitter", splitter)
pipeline.add_component("embedder", embedder)
pipeline.add_component("bm25_retriever", bm25_retriever)
pipeline.add_component("embedding_retriever", embedding_retriever)
pipeline.add_component("joiner", joiner)
pipeline.add_component("generator", generator)
# Connect components
pipeline.connect("cleaner", "splitter")
pipeline.connect("splitter", "embedder")
pipeline.connect("embedder", "embedding_retriever")
pipeline.connect("bm25_retriever", "joiner")
pipeline.connect("embedding_retriever", "joiner")
pipeline.connect("joiner", "generator")
# Add documents to store
from haystack.dataclasses import Document
documents = [
Document(content="The system requires 8GB RAM minimum"),
Document(content="Supports Windows, macOS, and Linux"),
Document(content="Network bandwidth should be at least 10Mbps")
]
document_store.write_documents(documents)
# Run pipeline
result = pipeline.run({
"cleaner": {"documents": documents},
"bm25_retriever": {"query": "system requirements"},
"embedding_retriever": {"query": "system requirements"}
})
print(result["generator"]["replies"][0])
混合搜索实现方式
from haystack import Pipeline
from haystack.components.embedders import SentenceTransformersTextEmbedder
from haystack.components.retrievers import InMemoryBM25Retriever, InMemoryEmbeddingRetriever
from haystack.components.rankers import TransformersRanker
from haystack.components.joiners import JoinDocuments
from haystack.components.generators import OpenAIGenerator
from haystack.document_stores.in_memory import InMemoryDocumentStore
# Create hybrid search pipeline
pipeline = Pipeline()
# Add components
document_store = InMemoryDocumentStore()
bm25_retriever = InMemoryBM25Retriever(document_store=document_store)
embedding_retriever = InMemoryEmbeddingRetriever(document_store=document_store)
ranker = TransformersRanker(model_name_or_path="BAAI/bge-reranker-base")
joiner = JoinDocuments(join_mode="concatenate")
generator = OpenAIGenerator(model="gpt-4o")
# Add components to pipeline
pipeline.add_component("bm25_retriever", bm25_retriever)
pipeline.add_component("embedding_retriever", embedding_retriever)
pipeline.add_component("ranker", ranker)
pipeline.add_component("joiner", joiner)
pipeline.add_component("generator", generator)
# Connect components
pipeline.connect("bm25_retriever", "ranker.query")
pipeline.connect("embedding_retriever", "ranker.documents")
pipeline.connect("ranker", "joiner")
pipeline.connect("joiner", "generator")
# Run hybrid search
result = pipeline.run({
"bm25_retriever": {"query": "system requirements"},
"embedding_retriever": {"query": "system requirements"}
})
print(result["generator"]["replies"][0])
5. RAGAS:评估框架
基本评估配置
import pandas as pd
from ragas import evaluate
from ragas.metrics import (
faithfulness,
answer_relevancy,
context_precision,
context_recall,
context_relevancy,
answer_similarity
)
from datasets import Dataset
# Create evaluation dataset
data = {
"question": ["What are the system requirements?", "How does the authentication work?"],
"answer": ["The system requires 8GB RAM minimum", "Authentication uses OAuth 2.0"],
"contexts": [
["The system requires 8GB RAM minimum", "Supports Windows, macOS, and Linux"],
["Authentication uses OAuth 2.0", "Multi-factor authentication is optional"]
],
"ground_truths": [
["The system requires 8GB RAM minimum"],
["Authentication uses OAuth 2.0 with JWT tokens"]
]
}
dataset = Dataset.from_dict(data)
# Evaluate
result = evaluate(
dataset,
metrics=[
faithfulness,
answer_relevancy,
context_precision,
context_recall,
context_relevancy,
answer_similarity
]
)
print(result)
使用无参考指标的自定义评估方法
from ragas import evaluate
from ragas.metrics import (
faithfulness,
answer_relevancy,
context_precision,
context_recall,
context_relevancy,
answer_similarity,
answer_correctness
)
from datasets import Dataset
# Create dataset without ground truth (reference-free evaluation)
data = {
"question": ["What are the system requirements?", "How does the authentication work?"],
"answer": ["The system requires 8GB RAM minimum", "Authentication uses OAuth 2.0"],
"contexts": [
["The system requires 8GB RAM minimum", "Supports Windows, macOS, and Linux"],
["Authentication uses OAuth 2.0", "Multi-factor authentication is optional"]
]
}
dataset = Dataset.from_dict(data)
# Evaluate without ground truth
result = evaluate(
dataset,
metrics=[
faithfulness,
answer_relevancy,
context_precision,
context_recall,
context_relevancy,
answer_similarity
]
)
print(result)
与现有RAG系统的集成
from ragas import evaluate
from ragas.metrics import faithfulness, answer_relevancy, context_precision
from datasets import Dataset
import json
def evaluate_rag_system(rag_system, test_questions, expected_answers=None):
"""Evaluate a RAG system against test questions"""
results = {
"question": [],
"answer": [],
"contexts": [],
"ground_truths": [] if expected_answers else None
}
for i, question in enumerate(test_questions):
# Get answer from RAG system
answer = rag_system(question)
# Get contexts used by RAG system (assuming it returns contexts)
contexts = rag_system.get_contexts(question) if hasattr(rag_system, 'get_contexts') else []
results["question"].append(question)
results["answer"].append(answer)
results["contexts"].append(contexts)
if expected_answers:
results["ground_truths"].append([expected_answers[i]])
# Create dataset
dataset = Dataset.from_dict(results)
# Evaluate
metrics = [faithfulness, answer_relevancy, context_precision]
if expected_answers:
metrics.append(context_recall)
evaluation_result = evaluate(dataset, metrics=metrics)
return evaluation_result
# Example usage
# Assuming you have a RAG system object
# evaluation = evaluate_rag_system(your_rag_system, test_questions, expected_answers)
# print(evaluation)
对比分析:在生产架构中整合多种工具
推荐的架构模式
# Complete production architecture combining all tools
from docling.document_converter import DocumentConverter
from llama_index.core import VectorStoreIndex, SimpleDirectoryReader
from langchain_core.prompts import ChatPromptTemplate
from haystack import Pipeline
from ragas import evaluate
from datasets import Dataset
class ProductionRAGSystem:
def __init__(self):
self.docling_converter = DocumentConverter()
self.llamaindex_index = None
self.langchain_chain = None
self.haystack_pipeline = None
self.evaluation_metrics = []
def preprocess_documents(self, document_paths):
"""Use Docling for intelligent document processing"""
processed_docs = []
for path in document_paths:
result = self.docling_converter.convert(path)
markdown_content = result.document.export_to_markdown()
processed_docs.append({
"content": markdown_content,
"metadata": {"source": path}
})
return processed_docs
def build_llamaindex_index(self, documents):
"""Build index using LlamaIndex for efficient retrieval"""
from llama_index.core import VectorStoreIndex
from llama_index.core.node_parser import SentenceSplitter
# Split documents
splitter = SentenceSplitter(chunk_size=512, chunk_overlap=50)
nodes = splitter.get_nodes_from_documents(documents)
# Create index
self.llamaindex_index = VectorStoreIndex(nodes)
return self.llamaindex_index
def create_langchain_chain(self, index):
"""Create LangChain chain for complex workflows"""
from langchain_core.runnables import RunnablePassthrough
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
# Create retriever
retriever = index.as_retriever(similarity_top_k=3)
# Create prompt
template = """Answer the question based only on the following context:
{context}
Question: {question}
"""
prompt = ChatPromptTemplate.from_template(template)
# Create chain
llm = ChatOpenAI(model_name="gpt-4o", temperature=0)
self.langchain_chain = (
{"context": retriever, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
return self.langchain_chain
def setup_haystack_pipeline(self, documents):
"""Set up Haystack pipeline for production deployment"""
from haystack import Pipeline
from haystack.components.embedders import SentenceTransformersTextEmbedder
from haystack.components.retrievers import InMemoryEmbeddingRetriever
from haystack.components.generators import OpenAIGenerator
from haystack.document_stores.in_memory import InMemoryDocumentStore
# Initialize components
document_store = InMemoryDocumentStore()
embedder = SentenceTransformersTextEmbedder(model="sentence-transformers/all-MiniLM-L6-v2")
retriever = InMemoryEmbeddingRetriever(document_store=document_store)
generator = OpenAIGenerator(model="gpt-4o")
# Create pipeline
pipeline = Pipeline()
pipeline.add_component("embedder", embedder)
pipeline.add_component("retriever", retriever)
pipeline.add_component("generator", generator)
# Connect components
pipeline.connect("embedder", "retriever")
pipeline.connect("retriever", "generator")
# Add documents
document_store.write_documents(documents)
self.haystack_pipeline = pipeline
return pipeline
def evaluate_system(self, test_questions, test_answers=None):
"""Evaluate system using RAGAS"""
data = {
"question": test_questions,
"answer": [],
"contexts": []
}
if test_answers:
data["ground_truths"] = []
# Generate answers
for question in test_questions:
answer = self.langchain_chain.invoke(question)
contexts = self.llamaindex_index.as_retriever().invoke(question)
data["answer"].append(answer)
data["contexts"].append([ctx.text for ctx in contexts])
if test_answers:
data["ground_truths"].append([test_answers[test_questions.index(question)]])
# Create dataset
dataset = Dataset.from_dict(data)
# Evaluate
from ragas.metrics import faithfulness, answer_relevancy, context_precision
metrics = [faithfulness, answer_relevancy, context_precision]
if test_answers:
from ragas.metrics import context_recall
metrics.append(context_recall)
result = evaluate(dataset, metrics=metrics)
self.evaluation_metrics.append(result)
return result
# Example usage
rag_system = ProductionRAGSystem()
# 1. Preprocess documents with Docling
processed_docs = rag_system.preprocess_documents(["doc1.pdf", "doc2.pdf"])
# 2. Build index with LlamaIndex
index = rag_system.build_llamaindex_index(processed_docs)
# 3. Create LangChain chain
chain = rag_system.create_langchain_chain(index)
# 4. Set up Haystack pipeline for production
haystack_pipeline = rag_system.setup_haystack_pipeline(processed_docs)
# 5. Evaluate system
test_questions = ["What are the system requirements?", "How does authentication work?"]
test_answers = ["The system requires 8GB RAM minimum", "Authentication uses OAuth 2.0"]
evaluation = rag_system.evaluate_system(test_questions, test_answers)
print(evaluation)
实现时的注意事项与最佳实践
1. 文档处理的最佳实践
# Best practices for document processing with Docling
from docling.document_converter import DocumentConverter
from docling.datamodel.pipeline_options import PipelineOptions
def optimize_docling_processing():
"""Optimize Docling for different document types"""
# For scanned documents with OCR
pipeline_options_ocr = PipelineOptions(
do_ocr=True,
ocr_engine="tesseract",
ocr_language="eng"
)
# For clean digital documents
pipeline_options_clean = PipelineOptions(
do_ocr=False,
extract_images=False # Disable image extraction for text-only processing
)
# For documents with complex layouts
pipeline_options_layout = PipelineOptions(
do_ocr=True,
ocr_engine="easyocr", # More accurate but slower
extract_tables=True,
extract_images=True
)
return {
"ocr": pipeline_options_ocr,
"clean": pipeline_options_clean,
"layout": pipeline_options_layout
}
# Usage
optimizations = optimize_docling_processing()
converter = DocumentConverter(pipeline_options=optimizations["ocr"])
2. 分块策略的优化
# Advanced chunking strategies for different content types
from llama_index.core.node_parser import (
SentenceSplitter,
SemanticSplitterNodeParser,
HierarchicalNodeParser
)
def create_optimal_chunking_strategy(content_type="general"):
"""Create optimal chunking strategy based on content type"""
if content_type == "technical":
# Technical documents need smaller chunks for precision
return SentenceSplitter(
chunk_size=512,
chunk_overlap=64,
paragraph_separator="\n\n",
sentence_separator="\\n"
)
elif content_type == "legal":
# Legal documents need to preserve entire clauses
return SemanticSplitterNodeParser(
buffer_size=1,
breakpoint_percentile_threshold=95,
embed_model="sentence-transformers/all-MiniLM-L6-v2"
)
elif content_type == "research":
# Research papers benefit from hierarchical chunking
return HierarchicalNodeParser.from_defaults(
chunk_sizes=[2048, 512, 128]
)
else:
# General purpose chunking
return SentenceSplitter(
chunk_size=1024,
chunk_overlap=128
)
# Usage
chunking_strategy = create_optimal_chunking_strategy("technical")
3. 生产部署注意事项
# Production deployment configuration for Haystack
from haystack import Pipeline
from haystack.components.embedders import SentenceTransformersTextEmbedder
from haystack.components.retrievers import InMemoryEmbeddingRetriever
from haystack.components.generators import OpenAIGenerator
from haystack.document_stores.in_memory import InMemoryDocumentStore
import os
def create_production_pipeline():
"""Create production-ready Haystack pipeline"""
# Use environment variables for sensitive data
api_key = os.getenv("OPENAI_API_KEY")
model_name = os.getenv("LLM_MODEL_NAME", "gpt-4o")
# Initialize components with production settings
document_store = InMemoryDocumentStore()
embedder = SentenceTransformersTextEmbedder(
model="sentence-transformers/all-MiniLM-L6-v2",
batch_size=32 # Optimize for throughput
)
retriever = InMemoryEmbeddingRetriever(
document_store=document_store,
top_k=5 # Return more results for ranking
)
generator = OpenAIGenerator(
api_key=api_key,
model=model_name,
max_tokens=1024,
temperature=0.7 # Balance creativity and accuracy
)
# Create pipeline
pipeline = Pipeline()
pipeline.add_component("embedder", embedder)
pipeline.add_component("retriever", retriever)
pipeline.add_component("generator", generator)
# Connect components
pipeline.connect("embedder", "retriever")
pipeline.connect("retriever", "generator")
return pipeline
# Add monitoring and logging
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
def monitored_pipeline_run(pipeline, query):
"""Run pipeline with monitoring"""
logger.info(f"Processing query: {query}")
try:
result = pipeline.run({"retriever": {"query": query}})
logger.info("Query processed successfully")
return result
except Exception as e:
logger.error(f"Error processing query: {e}")
raise
4. 评估与持续改进
# Continuous evaluation and improvement loop
from ragas import evaluate
from ragas.metrics import (
faithfulness,
answer_relevancy,
context_precision,
context_recall
)
from datasets import Dataset
import time
class RAGEvaluationSystem:
def __init__(self):
self.evaluation_history = []
self.improvement_plan = {}
def run_evaluation_cycle(self, rag_system, test_set, previous_results=None):
"""Run evaluation cycle and generate improvement plan"""
# Evaluate current system
evaluation_result = self.evaluate_system(rag_system, test_set)
# Compare with previous results
if previous_results:
improvement = self.calculate_improvement(evaluation_result, previous_results)
self.generate_improvement_plan(improvement, evaluation_result)
# Record results
self.evaluation_history.append({
"timestamp": time.time(),
"results": evaluation_result,
"improvement_plan": self.improvement_plan.copy()
})
return evaluation_result
def evaluate_system(self, rag_system, test_set):
"""Evaluate RAG system"""
data = {
"question": test_set["questions"],
"answer": [],
"contexts": [],
"ground_truths": test_set["answers"]
}
# Generate answers
for question in test_set["questions"]:
answer = rag_system(question)
contexts = rag_system.get_contexts(question) if hasattr(rag_system, 'get_contexts') else []
data["answer"].append(answer)
data["contexts"].append(contexts)
# Create dataset
dataset = Dataset.from_dict(data)
# Evaluate
metrics = [
faithfulness,
answer_relevancy,
context_precision,
context_recall
]
result = evaluate(dataset, metrics=metrics)
return result
def calculate_improvement(self, current, previous):
"""Calculate improvement between evaluation cycles"""
improvement = {}
for metric in current.keys():
if metric in previous:
improvement[metric] = current[metric] - previous[metric]
return improvement
def generate_improvement_plan(self, improvement, current_results):
"""Generate improvement plan based on evaluation results"""
# Identify areas needing improvement
if improvement.get("faithfulness", 0) < 0.1:
self.improvement_plan["faithfulness"] = "Improve retrieval quality by adjusting chunk size or using better embedding model"
if improvement.get("answer_relevancy", 0) < 0.1:
self.improvement_plan["answer_relevancy"] = "Improve prompt engineering or use more sophisticated answer generation techniques"
if improvement.get("context_precision", 0) < 0.1:
self.improvement_plan["context_precision"] = "Implement re-ranking or hybrid search to improve context selection"
if improvement.get("context_recall", 0) < 0.1:
self.improvement_plan["context_recall"] = "Increase top-k parameter or implement query expansion techniques"
# Example usage
evaluator = RAGEvaluationSystem()
# Define test set
test_set = {
"questions": ["What are the system requirements?", "How does authentication work?"],
"answers": [["The system requires 8GB RAM minimum"], ["Authentication uses OAuth 2.0"]]
}
# Run evaluation
evaluation_result = evaluator.run_evaluation_cycle(your_rag_system, test_set)
print(evaluation_result)
print("Improvement Plan:", evaluator.improvement_plan)