LC

NutriCore RAG Pipeline — LangChain v0.3+ LCEL

Pipeline RAG production-ready para documentacion tecnica interna · CULTIVA IA

✓ Produccion Claude Sonnet Chroma · MMR
nutricore_rag.py
IMPORTS & SETUP
1# nutricore_rag.py — Pipeline RAG produccion · NutriCore Analytics
2# LangChain v0.3+ LCEL · Claude Sonnet + GPT-4o fallback · Chroma MMR
3
4from langchain_community.document_loaders import (
5    DirectoryLoader, PyPDFLoader, TextLoader
6)
7from langchain_text_splitters import RecursiveCharacterTextSplitter
8from langchain_chroma import Chroma
9from langchain_openai import OpenAIEmbeddings, ChatOpenAI
10from langchain_anthropic import ChatAnthropic
11from langchain_core.prompts import ChatPromptTemplate
12from langchain_core.output_parsers import StrOutputParser
13from langchain_core.runnables import RunnablePassthrough, RunnableParallel
14from langchain_core.globals import set_llm_cache
15from langchain_community.cache import SQLiteCache
16from pydantic import BaseModel, Field
17from typing import List
18
STRUCTURED OUTPUT — RESPUESTA CON CITAS
19class Source(BaseModel):
20    doc_name: str = Field(description="Nombre del documento fuente")
21    page: int = Field(description="Numero de pagina (0 si no aplica)")
22    relevance: float = Field(description="Score 0-1 de relevancia")
23
24class RagAnswer(BaseModel):
25    answer: str = Field(description="Respuesta fundamentada en el contexto")
26    sources: List[Source] = Field(description="Lista de fuentes usadas")
27    confidence: float = Field(description="Confianza 0-1 de la respuesta")
28    needs_escalation: bool = Field(description="True si la pregunta necesita un humano")
29
LLM PRINCIPAL + FALLBACK + CACHE
30# Cache SQLite: evita llamadas LLM duplicadas en prod
31set_llm_cache(SQLiteCache(database_path=".nutricore_cache.db"))
32
33# LLM principal: Claude Sonnet 4 (Anthropic)
34llm_primary = ChatAnthropic(
35    model="claude-sonnet-4-20250514",
36    temperature=0,
37    max_tokens=2048
38)
39
40# Fallback automatico a GPT-4o si Anthropic falla
41llm_fallback = ChatOpenAI(model="gpt-4o", temperature=0)
42llm = llm_primary.with_fallbacks([llm_fallback])
43
44# LLM con structured output para citas tipadas
45llm_structured = llm_primary.with_structured_output(RagAnswer)
46
INGESTA DE DOCUMENTOS — CARGA + SPLIT + VECTOR STORE
47def build_vectorstore(docs_path: str = "./docs") -> Chroma:
48    """Carga documentos, los trocea y construye el vector store Chroma."""
49
50    # Carga PDFs + Markdown del directorio docs/
51    loaders = {
52        "**/*.pdf": PyPDFLoader,
53        "**/*.md": TextLoader,
54    }
55    all_docs = []
56    for glob_pattern, loader_cls in loaders.items():
57        loader = DirectoryLoader(
58            docs_path, glob=glob_pattern, loader_cls=loader_cls,
59            show_progress=True, use_multithreading=True
60        )
61        all_docs.extend(loader.load())
62
63    # RecursiveCharacterTextSplitter: split optimo para RAG
64    splitter = RecursiveCharacterTextSplitter(
65        chunk_size=1000,
66        chunk_overlap=200,   # 20% overlap para no romper contexto
67        separators=["\n\n", "\n", ". ", " ", ""]
68    )
69    splits = splitter.split_documents(all_docs)
70
71    # Embeddings OpenAI + Chroma persistido en disco
72    embeddings = OpenAIEmbeddings(model="text-embedding-3-large")
73    vectorstore = Chroma.from_documents(
74        splits, embeddings,
75        persist_directory="./chroma_db"
76    )
77    return vectorstore
78
RAG CHAIN — LCEL CON MMR RETRIEVAL + STRUCTURED OUTPUT
79def build_rag_chain(vectorstore: Chroma):
80    """Construye el pipeline RAG con LCEL. Retorna runnable con streaming."""
81
82    # MMR: Maximal Marginal Relevance — diversidad sobre pura similitud
83    retriever = vectorstore.as_retriever(
84        search_type="mmr",
85        search_kwargs={"k": 5, "fetch_k": 20}
86    )
87
88    # Prompt RAG con instruccion de citar fuentes
89    rag_prompt = ChatPromptTemplate.from_template("""
90Eres el asistente tecnico de NutriCore Analytics. Responde SOLO usando
91el contexto proporcionado. Si no encuentras la respuesta, indica que
92necesitas escalar a un humano (needs_escalation=true).
93
94CONTEXTO DE DOCUMENTOS:
95{context}
96
97PREGUNTA: {question}
98
99Responde con JSON estructurado incluyendo fuentes exactas (doc + pagina).
100""")
101
102    def format_docs_with_meta(docs):
103        return "\n\n---\n".join(
104            f"[{d.metadata.get('source','?')} p.{d.metadata.get('page',0)}]\n{d.page_content}"
105            for d in docs
106        )
107
108    # LCEL chain con RunnableParallel para pasar question y context
109    rag_chain = (
110        RunnableParallel({
111            "context": retriever | format_docs_with_meta,
112            "question": RunnablePassthrough()
113        })
114        | rag_prompt
115        | llm_structured   # → RagAnswer tipado
116    )
117    return rag_chain
118
STREAMING EN PRODUCCION + FASTAPI ENDPOINT
119from fastapi import FastAPI
120from fastapi.responses import StreamingResponse
121
122app = FastAPI(title="NutriCore Knowledge API")
123
124# Inicializacion lazy al startup
125rag_chain = None
126
127@app.on_event("startup")
128async def startup():
129    global rag_chain
130    vs = build_vectorstore()
131    rag_chain = build_rag_chain(vs)
132
133@app.post("/ask")
134async def ask(question: str) -> RagAnswer:
135    """Endpoint con structured output tipado."""
136    return await rag_chain.ainvoke(question)
137
138@app.get("/ask/stream")
139async def ask_stream(question: str):
140    """Streaming con Server-Sent Events para UI en tiempo real."""
141    stream_chain = (
142        RunnableParallel({"context": retriever | format_docs,
143                           "question": RunnablePassthrough()})
144        | rag_prompt | llm | StrOutputParser()
145    )
146    async def event_generator():
147        async for chunk in stream_chain.astream(question):
148            yield f"data: {chunk}\n\n"
149    return StreamingResponse(event_generator(), media_type="text/event-stream")
150
DEMO EN VIVO — NutriCore Knowledge Assistant
API Running :8000
LCEL PIPELINE — FLUJO DE DATOS
📂
DirectoryLoader
63 docs · 41 PDFs + 22 Markdown · docs/ recursivo multihilo
8.3s
✂️
RecursiveCharacterTextSplitter
chunk_size=1000 · overlap=200 · 847 chunks generados
0.4s
🔢
OpenAIEmbeddings (text-embedding-3-large)
847 vectores · dim=3072 · batch paralelo
12.1s
💾
Chroma persist_directory=./chroma_db
Index persistido en disco · reutilizable en reinicios
1.2s
🔍
MMR Retriever · k=5 fetch_k=20
Maximal Marginal Relevance · diversidad sobre similitud pura
~45ms
🧠
Claude Sonnet 4 → with_fallbacks([GPT-4o])
structured_output(RagAnswer) · cache SQLite activo
~1.8s
CONSULTA EJEMPLO — RESPUESTA TIPADA
NutriCore Knowledge Assistant
D
Diego (dev junior) · hace 2 min
¿Como funciona el sistema de autenticacion JWT en NutriCore? ¿Donde se configura el tiempo de expiracion de los tokens y que pasa cuando uno caduca?
NC
NutriCore Assistant · Claude Sonnet 4 · streaming
El sistema JWT de NutriCore usa RS256 (par de claves RSA) con tokens de acceso de 15 minutos y refresh tokens de 7 dias.

La configuracion vive en config/auth.py:

ACCESS_TOKEN_EXPIRE = timedelta(minutes=15)
REFRESH_TOKEN_EXPIRE = timedelta(days=7)
ALGORITHM = "RS256"

Cuando un access token caduca, el cliente hace POST /auth/refresh con el refresh token. Si el refresh tambien caduco, el usuario debe re-autenticarse. El endpoint de refresh rota el refresh token (one-time use) para evitar token replay attacks.

Fuentes verificadas: 📄 auth-service-design.pdf p.8 📄 api-security-guide.pdf p.3 📝 ARCHITECTURE.md
Confianza: 0.94 · needs_escalation: false · 3 fuentes
METRICAS DE PRODUCCION
63
Documentos indexados
↑ 41 PDF + 22 MD
847
Chunks en Chroma
chunk=1000 · overlap=200
1.8s
Latencia P50
↓ cache hit 0.2s
0.94
Confianza media
sin escalado en 87% preguntas
€0.003
Coste por consulta
Claude Sonnet · ~1.2k tokens
100%
Uptime API
Railway · fallback GPT-4o activo
Stack completo: LangChain v0.3+ · LCEL pipe operator · RunnableParallel · MMR Retrieval · Chroma 0.5 · OpenAIEmbeddings text-embedding-3-large · Claude Sonnet 4 (principal) + GPT-4o (fallback) · Pydantic structured output · SQLiteCache · FastAPI StreamingResponse · Railway deployment