Inicio / Artículos / Notas prácticas: Mi sistema RAG omitió el 80 % de mis datos. Este único cambio lo solucionó.

Notas prácticas: Mi sistema RAG omitió el 80 % de mis datos. Este único cambio lo solucionó.

Guía paso a paso para utilizar las notas prácticas: Mi sistema RAG ignoró el 80 % de mis datos. Este único cambio lo solucionó.: contratos, verificaciones y espacios para código adicional para los equipos que implementan este patrón.

5233 palabras

Esta guía reconstruye el proceso desde las materias primas hasta un sistema funcional para: Mi sistema RAG perdió el 80 % de mis datos. Este único cambio lo solucionó.. El enfoque está en pasos operativos, verificaciones explícitas y código que se puede insertar directamente en un repositorio sin tener que adivinar su propósito. En la fase de visión general, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido sin tener que adivinar el estado oculto. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin necesidad de leer todo el sistema.

¿Qué cambió realmente con Gemini Embedding 2?

Al trabajar en la etapa de “¿Qué cambió realmente?”, anote primero el contrato: los datos requeridos, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Documente junto con ello el camino óptimo y el camino de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son mejoras posteriores. Registre el ID de la solicitud, el ID del modelo y la latencia en cada llamada. Sin ese registro, los errores intermitentes del proveedor parecen bugs de la aplicación.

El problema del espacio único de incrustación

Al trabajar en la etapa del Espacio de Incrustación Única, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Prefiera unidades pequeñas y probables sobre scripts extensos. Cuando un paso falla, el fallo debe apuntar a una única responsabilidad y no a un proceso complicado. Mida la capacidad de recuperación con un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente.

Qué realmente soporta el modelo

Al trabajar en la etapa de “¿Qué hace realmente el modelo?”, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Trate esta etapa como un contrato entre las entradas y las salidas validadas. Asigne nombres a los artefactos, defina comprobaciones de éxito y rechace las completaciones parciales silenciosas. Almacene en caché las instrucciones del sistema estables y los esquemas de las herramientas. Reenviar un preámbulo idéntico es una causa común de desperdicio de recursos. Al trabajar en la etapa de “¿Qué hace realmente el modelo?”, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el sistema.

La arquitectura, explicada en detalle

La arquitectura explicada en etapas funciona mejor cuando se trata como una superficie medible. Capture un caso exitoso, un caso de fallo y la nota de reversión antes de ampliar el alcance. Documente tanto el camino óptimo como el camino de recuperación juntos. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son ajustes realizados posteriormente. Separe la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

El pipeline de ingestión

La etapa del pipeline de ingestión funciona mejor cuando se trata como una superficie medible. Capture una transcripción ejemplar, un caso de fallo y la nota de reversión antes de ampliar el alcance. Prefiera unidades pequeñas y verificables en lugar de scripts extensos. Cuando falla un paso, el error debe apuntar a una única responsabilidad y no a todo el pipeline complicado. Separe la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

El pipeline de consultas

La etapa del Pipeline de Consultas funciona mejor cuando se trata como una superficie medible. Capture un transcripte ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Trate esta etapa como un contrato entre las entradas y las salidas validadas. Asigne nombres a los artefactos, defina comprobaciones de éxito y rechace completaciones parciales silenciosas. Separe la política de fragmentación de la política de recuperación. Cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad. La etapa del Pipeline de Consultas funciona mejor cuando se trata como una superficie medible. Capture un transcripte ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el sistema.

Por qué estos dos pipelines deben permanecer separados

En la etapa “Why These Two Pipelines”, defina las entradas, el responsable de cada paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido, sin tener que adivinar el estado oculto. Documente tanto la ruta óptima como la ruta de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son mejoras posteriores. Cite los pasajes que realmente sustentan la respuesta; sin citas, los operadores no podrán distinguir entre alucinaciones y lagunas en el indexado.

Configuración del entorno

En la fase de configuración del entorno, se deben definir las entradas, el responsable de cada paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido, sin tener que adivinar el estado oculto. Es preferible utilizar unidades pequeñas y verificables en lugar de scripts extensos. Cuando un paso falla, el error debe indicar una única responsabilidad y no un proceso complicado. Se deben citar los pasajes que realmente sustentan la respuesta. Sin citas, los operadores no pueden distinguir entre una alucinación y una laguna en el indexado.

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)

Construcción de la tubería de ingestión

En la etapa de Construcción del pipeline de ingestión, defina las entradas, el responsable de la tarea y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar la tarea a partir de un punto de control conocido sin tener que adivinar el estado oculto. Trate esta etapa como un contrato entre las entradas y los resultados validados. Asigne nombres a los artefactos, defina verificaciones de éxito y rechace las completaciones parciales silenciosas. Cite los pasajes que realmente sustentan la respuesta. Sin citaciones, los operadores no pueden distinguir entre alucinaciones y brechas en el indexado. En la etapa de Construcción del pipeline de ingestión, defina las entradas, el responsable de la tarea y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar la tarea a partir de un punto de control conocido sin tener que adivinar el estado oculto. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de secretos y las banderas de funcionalidad deben estar en un lugar que los operadores puedan auditar sin necesidad de leerlo.

todo el grafo.

Paso 1: El cliente de incrustación

Al trabajar en la etapa del Paso 1 de incrustación, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Documente junto con ello el camino óptimo y el camino de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no de mejoras posteriores. Mida el recuerdo en un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente.

# 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

Paso 2: El fragmentador de medios

Al trabajar en la etapa 2, “Los medios”, anote primero el contrato: los datos de entrada requeridos, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación ayuda a mantener honestos los cambios posteriores en el código. Prefiera unidades pequeñas y verificables a scripts extensos. Cuando un paso falla, el fallo debe apuntar a una única responsabilidad y no a un proceso complicado. Mida el rendimiento en un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente.

# 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

Paso 3: El almacén de vectores

Al trabajar en la etapa Paso 3: El Vector, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Trate esta etapa como un contrato entre las entradas y los resultados validados. Asigne nombres a los artefactos, defina comprobaciones de éxito y rechace las completaciones parciales silenciosas. Mida el rendimiento en un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente. Al trabajar en la etapa Paso 3: El Vector, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el código.

# 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()

Paso 4: El Ejecutor de Ingestión

La etapa 4, la de ingestión, funciona mejor cuando se trata como una superficie medible. Capture una transcripción exitosa, un caso de fallo y la nota de reversión antes de ampliar el alcance. Documente tanto el camino óptimo como el de recuperación juntos. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son mejoras posteriores. Separe la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

# 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")

Construyendo la tubería de consultas

La etapa de construcción del pipeline de consultas funciona mejor cuando se trata como una superficie medible. Capture un registro ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Prefiera unidades pequeñas y verificables en lugar de scripts extensos. Cuando falla un paso, el error debe apuntar a una única responsabilidad y no a un pipeline complicado. Separe la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

# 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)
    }

Mejorando la precisión de la recuperación

La etapa de Mejora de la Precisión en la Recuperación funciona mejor cuando se trata como una superficie medible. Capture una transcripción ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Trate esta etapa como un contrato entre las entradas y las salidas validadas. Asigne nombres a los artefactos, defina verificaciones de éxito y rechace completaciones parciales silenciosas. Separe la política de fragmentación de la política de recuperación. Cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad. La etapa de Mejora de la Precisión en la Recuperación funciona mejor cuando se trata como una superficie medible. Capture una transcripción ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el sistema.

1. Utilice siempre el tipo de tarea adecuado

En la etapa 1 “Use the Right”, defina las entradas, el responsable de cada paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido, sin tener que adivinar el estado oculto. Documente tanto la ruta óptima como la ruta de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son mejoras posteriores. Cite los pasajes que realmente sustentan la respuesta; sin citas, los operadores no podrán distinguir entre alucinaciones y lagunas en el indexado.

2. Agregue la expansión de consultas antes del incrustado

En la etapa 2 de Expansión de Consultas, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido sin tener que adivinar el estado oculto. Prefiera unidades pequeñas y verificables en lugar de scripts extensos. Cuando un paso falla, el error debe indicar una única responsabilidad y no un proceso complicado. Cite los pasajes que realmente sustentan la respuesta. Sin citaciones, los operadores no pueden distinguir entre alucinaciones y brechas en el indexado.

# 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. Ejecutar múltiples consultas en paralelo

En la etapa de 3 Ejecutar múltiples consultas, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido sin tener que adivinar el estado oculto. Trate esta etapa como un contrato entre las entradas y los resultados validados. Asigne nombres a los artefactos, defina verificaciones de éxito y rechace las completaciones parciales silenciosas. Cite los pasajes que realmente sustentan la respuesta. Sin citaciones, los operadores no pueden distinguir entre alucinaciones y brechas en el indexado. En la etapa de 3 Ejecutar múltiples consultas, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido sin tener que adivinar el estado oculto. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos secretos y las banderas de funcionalidad deben estar en un lugar que los operadores puedan auditar sin tener que leer todo el sistema.

# 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. Reclasificar los fragmentos recuperados

Al trabajar en la fase de 4 Reclasificar los fragmentos recuperados, anote primero el contrato: las entradas requeridas, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Documente tanto el camino óptimo como el de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no de mejoras posteriores. Mida el recuerdo en un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente.

# 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. Establecer un umbral de similitud y ceñirse a él

Al trabajar en la etapa 5 de Establecer una similitud, anote primero el contrato: los datos necesarios, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación mantiene honestas las futuras modificaciones del código. Prefiera unidades pequeñas y probables sobre scripts extensos. Cuando un paso falla, el fallo debe apuntar a una única responsabilidad y no a un proceso complicado. Mida la capacidad de recuperación con un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente.

Los resultados

Al trabajar en la etapa de Resultados, anote primero el contrato: los datos de entrada requeridos, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación garantiza que los cambios posteriores en el código sean transparentes. Considere esta etapa como un contrato entre los datos de entrada y los resultados validados. Asigne nombres a los artefactos, defina las comprobaciones de éxito y evite completar tareas parcialmente sin notificarlo. Mida el rendimiento en un conjunto fijo de preguntas antes de ajustar los prompts. El cambio constante de prompts rara vez soluciona un sistema de recuperación deficiente. Al trabajar en la etapa de Resultados, anote primero el contrato: los datos de entrada requeridos, la señal de éxito y qué ocurre en caso de fallo parcial. Esa lista de verificación garantiza que los cambios posteriores en el código sean transparentes. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el código.

Qué significan las dimensiones de Matryoshka para su infraestructura

La etapa de “What the Matryoshka Dimensions” funciona mejor cuando se trata como una superficie medible. Capture un registro exitoso, un caso de fallo y la nota de reversión antes de ampliar el alcance. Documente tanto el camino óptimo como el camino de recuperación juntos. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son ajustes realizados posteriormente. Separe la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

Qué puede construir con esto

La etapa “Lo que puedes construir” funciona mejor cuando se trata como una superficie medible. Consigue un registro ejemplar, un caso de fallo y la nota de reversión antes de ampliar el alcance. Prefiere unidades pequeñas y verificables en lugar de scripts extensos. Cuando un paso falla, el error debe apuntar a una sola responsabilidad y no a un proceso complicado. Separa la política de fragmentación de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

Consejos prácticos

La etapa de Conclusiones Accionables funciona mejor cuando se trata como una superficie medible. Capture un registro ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Trate esta etapa como un contrato entre las entradas y las salidas validadas. Asigne nombres a los artefactos, defina verificaciones de éxito y rechace completaciones parciales silenciosas. Separe la política de fragmentación de la política de recuperación. Cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad. La etapa de Conclusiones Accionables funciona mejor cuando se trata como una superficie medible. Capture un registro ideal, un caso de fallo y la nota de reversión antes de ampliar el alcance. Mantenga la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el sistema.

¿Qué falta aún?

En la fase de “¿Qué falta aún?”, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido, sin tener que adivinar el estado oculto. Documente tanto la ruta óptima como la ruta de recuperación. Las reintentos, los controles humanos y el manejo de mensajes no entregados forman parte del producto, no son mejoras posteriores. Cite los pasajes que realmente sustentan la respuesta. Sin citas, los operadores no pueden distinguir entre alucinaciones y brechas en el indexado.

Continuemos aprendiendo juntos

En la fase de “Let’s Keep Learning”, defina las entradas, el responsable del paso y los criterios de finalización antes de modificar el código. Los operadores deben poder volver a ejecutar el paso a partir de un punto de control conocido, sin tener que adivinar el estado oculto. Es preferible utilizar unidades pequeñas y verificables en lugar de scripts extensos. Cuando un paso falla, el error debe indicar una única responsabilidad y no un proceso complicado. Cite los pasajes que realmente sustentan la respuesta; sin citas, los operadores no pueden distinguir entre alucinaciones y fallos en el indexado.

Lista de verificación operativa

La fase de la lista de verificación operativa funciona mejor cuando se trata como una superficie medible. Capture una transcripción clave, un caso de fallo y la nota de reversión antes de ampliar el alcance. Registre los tiempos y el costo en tokens o consultas junto con los resultados funcionales. Tener visibilidad del costo desde el principio evita facturas inesperadas al pasar del entorno de demostración a los entornos compartidos.

Separar la política de fragmentación de la política de recuperación. Cambiar una no debería obligar a reescribir la otra cuando cambian las métricas de calidad.

Añadir una prueba de funcionamiento que ejerza la ruta crítica en el proceso de integración continua utilizando configuraciones fijas, y no APIs pagadas en tiempo real, siempre que lo permitan los presupuestos.

Mantener la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos secretos y las banderas de funcionalidad deben estar en un lugar donde los operadores puedan auditarlos sin tener que leer todo el sistema.

Separar la política de fragmentación de la política de recuperación. Cambiar una no debería obligar a reescribir la otra cuando cambian las métricas de calidad.

Antes de promocionar la solución, congelar las versiones, capturar una transcripción de referencia para la ruta crítica y confirmar los pasos para realizar un rollback. Los entornos compartidos necesitan límites de velocidad, verificaciones de asignación y un responsable claro para la rotación de datos secretos. Es mejor priorizar una fiabilidad sólida que demostraciones ingeniosas pero puntuales.

Nota por lotes para b3567bc23c05: mantener las claves del proveedor fuera del repositorio, establecer un límite para los tokens por sesión y almacenar las transcripciones junto a los archivos de evaluación para que los cambios posteriores en el modelo sigan siendo comparables.