实用提示:我的RAG系统漏掉了80%的数据,这一项改动解决了问题。
《实用笔记》操作指南:我的RAG系统漏掉了80%的数据,这一项改动解决了问题:适用于采用该模式的团队的合同、检查清单以及可直接插入的代码模块。
本指南将逐步构建从原始材料到可运行系统的完整流程,适用于文章:我的RAG系统漏掉了80%的数据,这一改动解决了问题。重点在于可操作的步骤、明确的检查点以及可直接放入代码库的代码,无需猜测其用途。 在概览阶段,应在修改代码之前明确输入参数、各步骤的负责人以及完成标准。操作人员应能够从已知的检查点重新运行相应步骤,而无需推测隐藏的状态。 配置信息应与应用程序代码分开存放。环境文件、密钥存储以及功能开关应集中管理,以便操作人员无需查看整个系统结构即可进行审核。
Gemini Embedding 2究竟带来了哪些变化?
在处理“实际发生了什么变化”这一阶段时,首先需列出相关契约:所需的输入参数、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持一致性。 同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及死信处理都是产品本身的组成部分,而非后续的优化工作。 每次调用时都要记录请求ID、模型ID以及延迟时间。如果没有这些记录,间歇性的服务端错误就会被视为应用程序的缺陷。
单一嵌入空间问题
在处理“单一嵌入空间”阶段时,首先写下相关规范:所需的输入参数、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持一致性。 优先选择小型、可测试的单元,而非庞大的脚本。当某个步骤失败时,故障应能指向单一的责任模块,而非复杂的流程链。 在调整提示词之前,先使用固定的问题集来衡量召回率。仅仅更换提示词很难改善较差的检索效果。
模型实际支持的功能
在处理“模型实际功能”这一阶段时,首先需写下契约:所需输入、成功信号以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持透明。 将这一阶段视为输入与验证后输出之间的契约。为相关成果命名,明确成功判定标准,杜绝无声的半完成状态。 缓存稳定的系统指令和工具架构。重复发送相同的开头信息是导致资源浪费的常见原因。 在处理“模型实际功能”这一阶段时,首先需写下契约:所需输入、成功信号以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持透明。 将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一处,以便操作人员无需查看整个系统结构即可进行审计。
完整的架构说明
将“分阶段架构说明”视为可度量的对象来处理效果最佳。在扩大范围之前,先记录一个成功的案例、一个失败案例以及回滚说明。同时记录正常流程和恢复流程。重试机制、人工审核环节以及死信处理都是产品本身的一部分,而非后续需要补充的内容。应将分块策略与检索策略分开,当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
数据摄取流程
将“数据摄取管道”阶段视为可度量的对象来处理,其效果最佳。在扩大范围之前,先记录一份理想的转录结果、一个故障案例以及回滚说明。相比庞大的脚本,应优先使用小型且可测试的单元。当某个步骤出现故障时,故障原因应能明确指向某个具体的责任模块,而非整个复杂的管道系统。此外,应将分块策略与检索策略分开处理;当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
查询管道
将查询处理流程视为可度量的环节时,其效果最佳。在扩大范围之前,需记录一份理想的处理结果、一个失败案例以及回滚说明。 应将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,明确成功标准,杜绝默许的半完成状态。 需将分块策略与检索策略分开。当质量指标发生变化时,调整其中一项不应强制要求重新编写另一项。 将查询处理流程视为可度量的环节时,其效果最佳。在扩大范围之前,需记录一份理想的处理结果、一个失败案例以及回滚说明。 应将配置置于应用程序代码之外。环境文件、密钥存储及功能开关应集中存放于一处,以便操作人员无需查看整个系统结构即可进行审计。
为何这两套流程必须保持分离
在“为何选择这两个管道阶段”这一环节中,应在修改代码之前明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都属于产品功能的一部分,而非后续需要补充的内容。 必须引用实际作为答案依据的段落。如果没有引用,操作人员就无法区分是虚假信息还是索引缺失导致的问题。
环境配置
在“环境搭建”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏的状态。相比冗长的脚本,更应采用小型且可测试的单元。当某个步骤失败时,故障原因应能指向单一责任模块,而非复杂的流程链。需引用实际作为答案依据的段落;没有引用的话,操作人员就无法区分是虚假信息还是索引缺失所致。
pip install google-genai chromadb google-generativeai python-dotenv ffmpeg-python
# config.py
import os
from google import genai
from google.genai import types
GEMINI_API_KEY = os.getenv("GEMINI_API_KEY")
EMBEDDING_MODEL = "gemini-embedding-2-preview"
GENERATION_MODEL = "gemini-2.5-pro"
# Output dimensionality options: 128, 256, 512, 768, 1024, 1536, 3072
# 1536 is the recommended default
EMBEDDING_DIMENSIONS = 1536
client = genai.Client(api_key=GEMINI_API_KEY)
构建数据摄取管道
在构建数据摄取管道阶段,应在修改代码之前明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,定义成功检测标准,并拒绝默许的半完成状态。 需引用实际作为答案依据的段落。没有引用的话,操作人员就无法区分是幻觉还是索引缺失导致的错误。 在构建数据摄取管道阶段,应在修改代码之前明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于操作人员可审计且无需阅读的内容中。
整个图表。步骤1:嵌入客户端
在处理“嵌入”阶段的步骤1时,首先列出相关规范:所需输入、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续代码修改的规范性。 同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及死信处理都是产品功能的一部分,而非后续需要优化的内容。 在调整提示词之前,先使用固定的问题集来衡量召回率。仅仅更换提示词很难解决检索效果不佳的问题。
# embedder.py
import time
from pathlib import Path
from google import genai
from google.genai import types
from config import client, EMBEDDING_MODEL, EMBEDDING_DIMENSIONS
def embed_text(text: str, task_type: str = "RETRIEVAL_DOCUMENT") -> list[float]:
"""Embed a plain text chunk."""
result = client.models.embed_content(
model=EMBEDDING_MODEL,
contents=text,
config=types.EmbedContentConfig(
task_type=task_type,
output_dimensionality=EMBEDDING_DIMENSIONS
)
)
return result.embeddings[0].values
def _wait_for_file(uploaded, max_wait: int = 300):
"""Poll until a File API upload is done processing."""
waited = 0
poll_interval = 5
while uploaded.state.name == "PROCESSING" and waited time.sleep(poll_interval)
waited += poll_interval
uploaded = client.files.get(name=uploaded.name)
if uploaded.state.name != "ACTIVE":
raise RuntimeError(
f"File never became ACTIVE. Final state: {uploaded.state.name}"
)
return uploaded
def embed_audio(audio_path: str) -> list[float]:
"""
Embed an audio file natively. No transcription step.
The model processes the audio signal directly and returns a
semantic embedding that captures speech content, tone, and
acoustic features. Max input: 80 seconds per file.
"""
uploaded = client.files.upload(path=str(audio_path))
uploaded = _wait_for_file(uploaded, max_wait=120)
result = client.models.embed_content(
model=EMBEDDING_MODEL,
contents=uploaded,
config=types.EmbedContentConfig(
task_type="RETRIEVAL_DOCUMENT",
output_dimensionality=EMBEDDING_DIMENSIONS
)
)
# Clean up: uploaded files count against your quota
client.files.delete(name=uploaded.name)
return result.embeddings[0].values
def embed_video(video_path: str) -> list[float]:
"""
Embed a video chunk natively. Gemini processes both the
audio track and visual frames together in one pass.
This is the key capability: visual demonstrations get captured
in the embedding alongside what is being said. Max input: 128 seconds.
"""
uploaded = client.files.upload(path=str(video_path))
uploaded = _wait_for_file(uploaded, max_wait=300)
result = client.models.embed_content(
model=EMBEDDING_MODEL,
contents=uploaded,
config=types.EmbedContentConfig(
task_type="RETRIEVAL_DOCUMENT",
output_dimensionality=EMBEDDING_DIMENSIONS
)
)
client.files.delete(name=uploaded.name)
return result.embeddings[0].values
def embed_with_context(text: str, image_bytes: bytes = None) -> list[float]:
"""
Embed text and an optional image together in a single call.
When both are passed, the model returns one vector that
represents the joint meaning. A query asking about a database
schema can retrieve a screenshot of that schema.
"""
contents = [text]
if image_bytes:
contents.append(
types.Part.from_bytes(data=image_bytes, mime_type="image/jpeg")
)
result = client.models.embed_content(
model=EMBEDDING_MODEL,
contents=contents,
config=types.EmbedContentConfig(
task_type="RETRIEVAL_DOCUMENT",
output_dimensionality=EMBEDDING_DIMENSIONS
)
)
return result.embeddings[0].values
步骤2:媒体分块器
在完成“第2步:媒体处理”阶段时,首先写下相关契约:所需的输入参数、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持一致性。 建议使用小型、可测试的单元而非庞大的脚本。当某个步骤失败时,故障应指向单一的责任模块,而非复杂的处理流程。 在调整提示词之前,先使用固定的问题集来测试召回率。仅仅更换提示词很难解决检索效果不佳的问题。
# chunker.py
import subprocess
import json
from pathlib import Path
from dataclasses import dataclass
from typing import List
@dataclass
class MediaChunk:
file_path: str
start_time: float
end_time: float
source_file: str
modality: str
chunk_index: int
total_chunks: int # Useful for progress reporting
def get_media_duration(file_path: str) -> float:
"""Get exact duration via ffprobe. Works for both audio and video."""
cmd = [
"ffprobe", "-v", "quiet",
"-print_format", "json",
"-show_streams", str(file_path)
]
result = subprocess.run(cmd, capture_output=True, text=True, check=True)
data = json.loads(result.stdout)
# Find the first stream with a duration value
for stream in data.get("streams", []):
if "duration" in stream:
return float(stream["duration"])
raise ValueError(f"Could not determine duration for: {file_path}")
def _run_ffmpeg_split(input_path: str, output_path: str,
start: float, duration: float):
"""Execute a single ffmpeg split operation."""
cmd = [
"ffmpeg", "-y",
"-ss", str(start),
"-i", str(input_path),
"-t", str(duration),
"-c", "copy", # No re-encoding: much faster, no quality loss
"-avoid_negative_ts", "make_zero",
str(output_path)
]
result = subprocess.run(cmd, capture_output=True)
if result.returncode != 0:
raise RuntimeError(
f"ffmpeg failed: {result.stderr.decode()}"
)
def chunk_video(
video_path: str,
chunk_duration: int = 90,
overlap: int = 10,
output_dir: str = "./chunks/video"
) -> List[MediaChunk]:
"""
Split video into overlapping chunks within the 128-second limit.
Default: 90-second chunks with 10-second overlap.
Overlap ensures topic transitions are captured in at least one chunk.
"""
Path(output_dir).mkdir(parents=True, exist_ok=True)
total_duration = get_media_duration(video_path)
source_name = Path(video_path).stem
# Pre-calculate chunk boundaries
boundaries = []
start = 0.0
while start end = min(start + chunk_duration, total_duration)
boundaries.append((start, end))
start += (chunk_duration - overlap)
chunks = []
for idx, (start, end) in enumerate(boundaries):
output_path = f"{output_dir}/{source_name}_{idx:04d}.mp4"
_run_ffmpeg_split(video_path, output_path, start, end - start)
chunks.append(MediaChunk(
file_path=output_path,
start_time=start,
end_time=end,
source_file=str(video_path),
modality="video",
chunk_index=idx,
total_chunks=len(boundaries)
))
return chunks
def chunk_audio(
audio_path: str,
chunk_duration: int = 60,
overlap: int = 5,
output_dir: str = "./chunks/audio"
) -> List[MediaChunk]:
"""
Split audio into overlapping chunks within the 80-second limit.
60 seconds per chunk gives a comfortable buffer under the 80-second cap.
"""
Path(output_dir).mkdir(parents=True, exist_ok=True)
total_duration = get_media_duration(audio_path)
source_name = Path(audio_path).stem
boundaries = []
start = 0.0
while start end = min(start + chunk_duration, total_duration)
boundaries.append((start, end))
start += (chunk_duration - overlap)
chunks = []
for idx, (start, end) in enumerate(boundaries):
output_path = f"{output_dir}/{source_name}_{idx:04d}.mp3"
_run_ffmpeg_split(audio_path, output_path, start, end - start)
chunks.append(MediaChunk(
file_path=output_path,
start_time=start,
end_time=end,
source_file=str(audio_path),
modality="audio",
chunk_index=idx,
total_chunks=len(boundaries)
))
return chunks
第3步:向量存储
在处理第3步“向量生成”阶段时,首先需写下相关契约:所需输入、成功信号以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合约定。 将此阶段视为输入与验证后输出之间的契约。为相关产物命名,明确成功判定标准,杜绝无声的半完成状态。 在调整提示词之前,先使用固定的问题集来测试召回率。仅仅更换提示词往往无法解决检索效果不佳的问题。 在处理第3步“向量生成”阶段时,首先需写下相关契约:所需输入、成功信号以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合约定。 将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。
# vector_store.py
import chromadb
from chromadb.config import Settings
from pathlib import Path
class MultimodalVectorStore:
"""
Vector store wrapping ChromaDB for multimodal RAG.
Stores embeddings + metadata for text, audio, and video chunks.
"""
def __init__(self, persist_dir: str = "./chroma_db"):
self.client = chromadb.PersistentClient(
path=persist_dir,
settings=Settings(anonymized_telemetry=False)
)
self.collection = self.client.get_or_create_collection(
name="multimodal_rag",
# cosine distance is standard for semantic similarity
metadata={"hnsw:space": "cosine"}
)
def add_text_chunk(
self,
chunk_id: str,
text: str,
embedding: list[float],
source_file: str,
chunk_index: int,
page: int = None
):
self.collection.add(
ids=[chunk_id],
embeddings=[embedding],
documents=[text],
metadatas=[{
"modality": "text",
"source_file": source_file,
"chunk_index": chunk_index,
"page": page or 0,
"preview": text[:250]
}]
)
def add_media_chunk(
self,
chunk_id: str,
embedding: list[float],
source_file: str,
start_time: float,
end_time: float,
modality: str,
chunk_index: int
):
"""
Store a video or audio chunk.
Note: we store a formatted timestamp string in `documents`
so ChromaDB has something to display. The actual retrieval
quality comes entirely from the embedding, not this text.
"""
ts_start = f"{int(start_time // 60):02d}:{int(start_time % 60):02d}"
ts_end = f"{int(end_time // 60):02d}:{int(end_time % 60):02d}"
display = (
f"[{modality.upper()}] {Path(source_file).name} "
f"from {ts_start} to {ts_end}"
)
self.collection.add(
ids=[chunk_id],
embeddings=[embedding],
documents=[display],
metadatas=[{
"modality": modality,
"source_file": source_file,
"start_time": start_time,
"end_time": end_time,
"timestamp_start": ts_start,
"timestamp_end": ts_end,
"chunk_index": chunk_index,
"preview": display
}]
)
def search(
self,
query_embedding: list[float],
n_results: int = 5,
modality_filter: str = None
) -> list[dict]:
"""
Retrieve top-k most similar chunks across all modalities.
Optionally filter to a single modality for targeted search.
"""
where_clause = {"modality": modality_filter} if modality_filter else None
results = self.collection.query(
query_embeddings=[query_embedding],
n_results=n_results,
where=where_clause,
include=["documents", "metadatas", "distances"]
)
chunks = []
for doc, meta, dist in zip(
results["documents"][0],
results["metadatas"][0],
results["distances"][0]
):
chunks.append({
"content": doc,
"metadata": meta,
"modality": meta["modality"],
# ChromaDB returns cosine distance; convert to similarity score
"similarity": round(1.0 - dist, 4)
})
return sorted(chunks, key=lambda x: x["similarity"], reverse=True)
def count(self) -> int:
return self.collection.count()
第4步:数据摄取运行器
第4步:在“数据摄取”阶段,若将其视为可度量的指标,则效果最佳。在扩大范围之前,先记录一个成功的处理案例、一个失败案例以及回滚说明。同时将正常流程和恢复流程都记录下来。重试机制、人工审核环节以及死信处理都是产品本身的一部分,而非后续需要补充的功能。应将分块策略与检索策略分开,当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
# ingest.py
import os
import hashlib
from pathlib import Path
from chunker import chunk_video, chunk_audio
from embedder import embed_text, embed_audio, embed_video
from vector_store import MultimodalVectorStore
store = MultimodalVectorStore(persist_dir="./chroma_db")
def make_chunk_id(source_path: str, chunk_index: int) -> str:
"""Stable, unique ID for any chunk. Same input always = same ID."""
raw = f"{os.path.abspath(source_path)}:{chunk_index}"
return hashlib.sha256(raw.encode()).hexdigest()[:20]
def ingest_text_file(file_path: str):
with open(file_path, "r", encoding="utf-8") as f:
text = f.read()
# Sliding window chunking: 800 chars with 100-char overlap
chunk_size, overlap = 800, 100
raw_chunks = []
start = 0
while start end = min(start + chunk_size, len(text))
raw_chunks.append(text[start:end])
start += chunk_size - overlap
for i, chunk_text in enumerate(raw_chunks):
embedding = embed_text(chunk_text, task_type="RETRIEVAL_DOCUMENT")
store.add_text_chunk(
chunk_id=make_chunk_id(file_path, i),
text=chunk_text,
embedding=embedding,
source_file=file_path,
chunk_index=i
)
print(f" Stored {len(raw_chunks)} text chunks from {Path(file_path).name}")
def ingest_video_file(file_path: str):
print(f" Chunking: {Path(file_path).name}")
chunks = chunk_video(file_path, chunk_duration=90, overlap=10)
for chunk in chunks:
print(
f" Embedding chunk {chunk.chunk_index + 1}/{chunk.total_chunks} "
f"({chunk.start_time:.0f}s to {chunk.end_time:.0f}s)"
)
try:
embedding = embed_video(chunk.file_path)
store.add_media_chunk(
chunk_id=make_chunk_id(file_path, chunk.chunk_index),
embedding=embedding,
source_file=file_path,
start_time=chunk.start_time,
end_time=chunk.end_time,
modality="video",
chunk_index=chunk.chunk_index
)
except Exception as e:
print(f" WARNING: Failed to embed chunk {chunk.chunk_index}: {e}")
finally:
# Always clean up temp files, even on failure
if os.path.exists(chunk.file_path):
os.remove(chunk.file_path)
print(f" Done. {len(chunks)} video chunks stored.")
def ingest_audio_file(file_path: str):
print(f" Chunking: {Path(file_path).name}")
chunks = chunk_audio(file_path, chunk_duration=60, overlap=5)
for chunk in chunks:
try:
embedding = embed_audio(chunk.file_path)
store.add_media_chunk(
chunk_id=make_chunk_id(file_path, chunk.chunk_index),
embedding=embedding,
source_file=file_path,
start_time=chunk.start_time,
end_time=chunk.end_time,
modality="audio",
chunk_index=chunk.chunk_index
)
except Exception as e:
print(f" WARNING: Failed to embed chunk {chunk.chunk_index}: {e}")
finally:
if os.path.exists(chunk.file_path):
os.remove(chunk.file_path)
print(f" Done. {len(chunks)} audio chunks stored.")
def ingest_directory(directory: str):
handlers = {
".txt": ingest_text_file,
".md": ingest_text_file,
".mp4": ingest_video_file,
".mov": ingest_video_file,
".mp3": ingest_audio_file,
".wav": ingest_audio_file,
}
all_files = list(Path(directory).rglob("*"))
media_files = [f for f in all_files if f.suffix.lower() in handlers]
print(f"Found {len(media_files)} files to ingest\n")
for file_path in media_files:
print(f"Processing: {file_path.name}")
handler = handlers[file_path.suffix.lower()]
handler(str(file_path))
print()
print(f"Ingestion complete. Total chunks indexed: {store.count()}")
if __name__ == "__main__":
ingest_directory("./knowledge_base")
构建查询管道
将“构建查询流水线”阶段视为可度量的对象来处理效果最佳。在扩大范围之前,先记录一份理想的处理结果、一个失败案例以及回滚说明。 优先选择小型且可测试的单元,而非庞大的脚本。当某个步骤出现故障时,故障应指向单一责任模块,而非复杂的流水线。 将分块策略与检索策略分开。当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
# query.py
import os
from pathlib import Path
import google.generativeai as genai
from embedder import embed_text
from vector_store import MultimodalVectorStore
from config import GENERATION_MODEL
store = MultimodalVectorStore(persist_dir="./chroma_db")
def format_context_for_llm(chunks: list[dict]) -> str:
"""
Format retrieved chunks into a context block for the generative model.
We include modality, source, and similarity score so the model
can calibrate its confidence and cite sources accurately.
"""
parts = []
for rank, chunk in enumerate(chunks, start=1):
meta = chunk["metadata"]
modality = chunk["modality"]
score = chunk["similarity"]
if modality == "text":
parts.append(
f"[SOURCE {rank} | TEXT | {Path(meta['source_file']).name} "
f"| chunk {meta['chunk_index']} | similarity {score}]\n"
f"{chunk['content']}"
)
elif modality == "video":
parts.append(
f"[SOURCE {rank} | VIDEO | {Path(meta['source_file']).name} "
f"| {meta['timestamp_start']} to {meta['timestamp_end']} "
f"| similarity {score}]\n"
f"Video segment covering this time range."
)
elif modality == "audio":
parts.append(
f"[SOURCE {rank} | AUDIO | {Path(meta['source_file']).name} "
f"| {meta['timestamp_start']} to {meta['timestamp_end']} "
f"| similarity {score}]\n"
f"Audio segment covering this time range."
)
return "\n\n---\n\n".join(parts)
def answer_query(
query: str,
n_results: int = 5,
modality_filter: str = None,
similarity_threshold: float = 0.6
) -> dict:
"""
Full RAG pipeline: embed the query, retrieve chunks, generate answer.
similarity_threshold: chunks below this score are dropped before generation.
Prevents low-quality matches from polluting the context.
"""
query_embedding = embed_text(query, task_type="RETRIEVAL_QUERY")
retrieved = store.search(
query_embedding=query_embedding,
n_results=n_results,
modality_filter=modality_filter
)
# Filter out weak matches
filtered = [c for c in retrieved if c["similarity"] >= similarity_threshold]
if not filtered:
return {
"answer": (
"No sufficiently relevant content was found in the knowledge base. "
"The most similar content had a similarity score below the threshold."
),
"sources": retrieved,
"query": query
}
context = format_context_for_llm(filtered)
system_prompt = """You are a helpful assistant with access to a multimodal
knowledge base that contains text documents, video recordings, and audio files.
When citing a source, reference it by its label (e.g., SOURCE 1, SOURCE 2).
For video and audio sources, always include the timestamp so the user can
navigate to the exact moment in the recording.
If the retrieved context does not contain enough information to answer
confidently, say so clearly rather than guessing."""
user_message = (
f"Using only the sources below, answer this question:\n\n"
f"Question: {query}\n\n"
f"Sources:\n{context}"
)
genai.configure(api_key=os.getenv("GEMINI_API_KEY"))
model = genai.GenerativeModel(GENERATION_MODEL)
response = model.generate_content(
user_message,
generation_config={"temperature": 0.1}
)
return {
"answer": response.text,
"sources": filtered,
"query": query,
"chunks_retrieved": len(retrieved),
"chunks_used": len(filtered)
}
提升检索精度
将“提升检索精度”阶段视为可度量的工作面时,其效果最佳。在扩大范围之前,需记录一份理想样本、一个失败案例以及回滚说明。 应将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,明确成功标准,杜绝默许的半完成状态。 应将分块策略与检索策略分开。当质量指标发生变化时,修改其中一项不应强制要求重新编写另一项。 将“提升检索精度”阶段视为可度量的工作面时,其效果最佳。在扩大范围之前,需记录一份理想样本、一个失败案例以及回滚说明。 应将配置置于应用程序代码之外。环境文件、密钥存储和功能开关应集中存放于一处,以便操作人员无需查看整个系统结构即可进行审计。
1. 每次都使用正确的任务类型
在“1 使用正确阶段”中,应在修改代码之前明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都是产品本身的组成部分,而非后续需要补充的内容。 必须引用那些真正作为答案依据的段落。如果没有引用,操作人员就无法区分是幻觉内容还是索引缺失导致的错误。
2. 在嵌入之前进行查询扩展
在“2 添加查询扩展”阶段,修改代码之前需明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 优先选择小型、可测试的单元,而非冗长的脚本。当某个步骤失败时,故障应指向单一责任点,而非复杂的流程链。 需引用实际作为答案依据的段落。没有引用的话,操作人员就无法区分是幻觉内容还是索引缺失所致。
# query_expander.py
import google.generativeai as genai
import os
genai.configure(api_key=os.getenv("GEMINI_API_KEY"))
def expand_query(raw_query: str) -> str:
"""
Rewrite a short user query into a more detailed retrieval query.
Returns the expanded version. Falls back to original on failure.
"""
model = genai.GenerativeModel("gemini-2.0-flash")
prompt = (
"Rewrite the following search query to be more detailed and specific. "
"Add relevant context, related terminology, and clarify the intent. "
"Keep it as a single question. Do not add facts not implied by the original.\n\n"
f"Original query: {raw_query}\n\n"
"Expanded query:"
)
try:
response = model.generate_content(
prompt,
generation_config={"temperature": 0.2, "max_output_tokens": 200}
)
return response.text.strip()
except Exception:
return raw_query # Graceful fallback
# Usage in query pipeline:
# expanded = expand_query("API limits engineering review")
# query_embedding = embed_text(expanded, task_type="RETRIEVAL_QUERY")
3. 并行运行多个查询
在“执行多次查询”阶段,修改代码之前需先明确输入参数、该步骤的负责人以及终止标准。操作员应能够从已知的检查点重新运行该步骤,而无需猜测其中的隐藏状态。 应将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,设定成功判定标准,并杜绝无声的半完成状态。 必须引用实际作为答案依据的段落。没有引用的话,操作员就无法区分是幻觉内容还是索引缺失导致的错误。 在“执行多次查询”阶段,修改代码之前需先明确输入参数、该步骤的负责人以及终止标准。操作员应能够从已知的检查点重新运行该步骤,而无需猜测其中的隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能开关都应存放于操作员能够审核的位置,无需阅读整个系统结构。
# multi_query.py
import concurrent.futures
from query_expander import expand_query
from embedder import embed_text
from vector_store import MultimodalVectorStore
store = MultimodalVectorStore(persist_dir="./chroma_db")
def generate_query_variants(query: str) -> list[str]:
"""Generate multiple phrasings for a single question."""
import google.generativeai as genai
import os
genai.configure(api_key=os.getenv("GEMINI_API_KEY"))
model = genai.GenerativeModel("gemini-2.0-flash")
prompt = (
f"Generate 3 different ways to search for information about: {query}\n\n"
"Return exactly 3 queries, one per line, no numbering or bullets."
)
response = model.generate_content(prompt)
variants = [line.strip() for line in response.text.strip().split("\n") if line.strip()]
return ([query] + variants)[:4] # Always include original, cap at 4 total
def multi_query_search(query: str, n_per_query: int = 4) -> list[dict]:
"""
Search with multiple query variants and deduplicate results.
Returns unique chunks ranked by their best similarity score.
"""
variants = generate_query_variants(query)
def search_one(variant: str) -> list[dict]:
embedding = embed_text(variant, task_type="RETRIEVAL_QUERY")
return store.search(embedding, n_results=n_per_query)
# Run all variants in parallel
all_results = []
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures = {executor.submit(search_one, v): v for v in variants}
for future in concurrent.futures.as_completed(futures):
all_results.extend(future.result())
# Deduplicate by source file + chunk index, keeping best similarity score
seen = {}
for chunk in all_results:
meta = chunk["metadata"]
key = f"{meta['source_file']}:{meta['chunk_index']}"
if key not in seen or chunk["similarity"] > seen[key]["similarity"]:
seen[key] = chunk
return sorted(seen.values(), key=lambda x: x["similarity"], reverse=True)
4. 重新排序检索到的片段
在处理“重新排序检索到的片段”这一阶段时,首先明确相关要求:所需输入、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改不会偏离原有设计。 同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都属于产品功能的一部分,而非后续的优化工作。 在调整提示词之前,先使用固定的问题集来衡量召回率。仅仅更换提示词很难解决检索效果不佳的问题。
# reranker.py
from sentence_transformers import CrossEncoder
# This model runs locally, no API cost, fast inference
_reranker = None
def get_reranker():
global _reranker
if _reranker is None:
_reranker = CrossEncoder("cross-encoder/ms-marco-MiniLM-L-6-v2")
return _reranker
def rerank_chunks(query: str, chunks: list[dict], top_k: int = 5) -> list[dict]:
"""
Re-rank retrieved chunks using a cross-encoder model.
Cross-encoders read both query and chunk together, giving much
more precise relevance scores than embedding cosine similarity.
Only practical on a small candidate set (10-20 chunks).
"""
reranker = get_reranker()
# For video/audio, we use the metadata preview as the text input.
# For text chunks, we use the actual content.
pairs = []
for chunk in chunks:
if chunk["modality"] == "text":
doc_text = chunk["content"]
else:
meta = chunk["metadata"]
doc_text = (
f"{chunk['modality']} recording: {meta['source_file']} "
f"at {meta.get('timestamp_start', '')} to {meta.get('timestamp_end', '')}"
)
pairs.append([query, doc_text])
scores = reranker.predict(pairs)
for chunk, score in zip(chunks, scores):
chunk["rerank_score"] = float(score)
return sorted(chunks, key=lambda x: x["rerank_score"], reverse=True)[:top_k]
# Usage:
# candidates = store.search(query_embedding, n_results=20) # Retrieve wide
# final = rerank_chunks(query, candidates, top_k=5) # Re-rank narro
5. 设定相似度阈值并严格遵守
在完成“设置相似度”这5个步骤时,首先需写下相关约定:所需的输入参数、成功标志,以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持一致性。 建议使用小型、可测试的单元,而非冗长的脚本。当某个步骤失败时,故障应指向单一责任模块,而非复杂的流程链。 在调整提示词之前,先使用固定的问题集来衡量召回率。仅仅更换提示词很难解决检索效果不佳的问题。
结果
在处理“结果生成”阶段时,首先需明确相关约定:所需的输入参数、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合约定。 将这一阶段视为输入与验证后输出之间的契约。为相关成果命名,定义成功判定标准,杜绝无声的半完成状态。 在调整提示词之前,先使用固定的问题集来测试召回率。仅仅更换提示词往往无法解决检索效果不佳的问题。 在处理“结果生成”阶段时,首先需明确相关约定:所需的输入参数、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合约定。 将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。
马特罗什卡维度对您基础设施的意义
将“马特里奥什卡维度”阶段视为可测量的界面时,其效果最佳。在扩大范围之前,需记录一个成功的案例、一个失败案例以及回滚说明。同时记录正常流程与恢复流程。重试机制、人工审核环节以及死信处理都是产品本身的组成部分,而非后续需要补充的内容。应将分块策略与检索策略分开,当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
利用此框架可构建什么
“你可以构建什么”这一阶段若被视为可度量的基准,效果会更好。在扩大范围之前,先记录一份优秀的实现案例、一个失败案例以及回滚说明。 优先选择小型且可测试的单元,而非庞大的脚本。当某个步骤出错时,故障应能指向单一责任模块,而非复杂的流程链。 将分块策略与检索策略分开。当质量指标发生变化时,修改其中一项不应迫使重新编写另一项。
实用要点
“可操作总结”阶段若被视为可度量的工作面,效果最佳。在扩大范围之前,需记录一份最优案例、一个失败案例以及回滚说明。 应将此阶段视为输入与已验证输出之间的契约。为相关成果命名,明确成功标准,杜绝默许的半完成状态。 需将分块策略与检索策略分开。当质量指标发生变化时,调整其中一项不应强制重新编写另一项。 “可操作总结”阶段若被视为可度量的工作面,效果最佳。在扩大范围之前,需记录一份最优案例、一个失败案例以及回滚说明。 应将配置置于应用程序代码之外。环境文件、密钥存储及功能开关应集中存放于一处,以便操作人员无需查看整个系统结构即可进行审计。
还缺少什么
在“还有哪些缺失”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及错误处理都是产品本身的组成部分,而非后续需要补充的内容。 要引用那些真正作为答案依据的段落。如果没有引用,操作人员就无法区分是虚假信息还是索引缺失导致的错误。
让我们一起持续学习
在“持续学习”阶段,应在修改代码之前明确输入内容、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。相比冗长的脚本,更应优先使用小型、可测试的单元。当某个步骤失败时,故障原因应能指向单一责任点,而非复杂的流程链。必须引用实际作为答案依据的段落;没有引用的话,操作人员就无法区分是幻觉内容还是索引缺失导致的错误。
运营检查清单
将“运营检查清单”阶段视为可衡量的指标时,其效果最佳。在扩大范围之前,需记录一份标准示例、一个故障案例以及回滚说明。除了功能结果外,还需记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。
将分块策略与检索策略分开。当质量指标发生变化时,修改其中一项不应强制重新编写另一项。
在预算允许的情况下,使用测试数据而非真实的付费 API,在持续集成过程中添加用于检测关键路径的冒烟测试。
将配置置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。
将分块策略与检索策略分开。当质量指标发生变化时,修改其中一项不应强制重新编写另一项。
在升级技术栈之前,先冻结版本,为关键路径生成标准参考记录,并确认回滚步骤。共享环境需要设置速率限制、租户验证机制,以及明确的密钥轮换负责人。与其追求华丽的临时演示,不如注重扎实的可靠性。
b3567bc23c05的批处理说明:不要将提供者密钥放入代码仓库,为每个会话设置令牌上限,并将转录内容存储在评估测试用例的旁边,以便后续更换模型时仍能保持可比性。