Inicio / Artículos / Su agente RAG olvida todo después de un mensaje: aquí está cómo lo solucionó con Databricks…

Su agente RAG olvida todo después de un mensaje: aquí está cómo lo solucionó con Databricks…

Guía paso a paso para solucionar el problema de que su agente RAG olvida todo después de un mensaje: contratos, verificaciones y espacios para código listos para usar en equipos.

3410 palabras

Esta guía reconstruye el proceso desde las materias primas hasta un sistema funcional para el caso en que tu agente RAG olvida todo después de un único mensaje: así es como lo solucioné con Databricks Lakebase. Se enfoca en pasos operativos, verificaciones explícitas y código que puedes incorporar directamente a un repositorio sin tener que adivinar la intención. Para obtener una visión general, define 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 necesidad de adivinar el estado oculto. Considera esta etapa como un contrato entre las entradas y los resultados validados. Nombra los artefactos, define las verificaciones de éxito y rechaza cualquier completación parcial silenciosa.

La arquitectura

Al trabajar en La Arquitectura, 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. Registre los tiempos y el costo de tokens o consultas junto con los resultados funcionales. Tener visibilidad del costo desde el principio evita facturas inesperadas cuando el proceso pasa de la fase de demostración a entornos compartidos. 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.

PDFs (text + images + diagrams)
    ↓ ai_parse_document() (Version 2.0)
Parsed elements (text, tables, figure descriptions)
    ↓ RecursiveCharacterTextSplitter
Chunks table (Delta, with Change Data Feed)
    ↓ Delta Sync + GTE-Large
Vector Search Index
    ↓ VectorSearchRetrieverTool
LangChain Agent + PostgresSaver
    ↓                    ↓
LLM (Foundation Model)   Lakebase (conversation memory)
    ↓
Model Serving Endpoint (MLflow ResponsesAgent)

Paso 1: Analizar documentos con ai_parse_document()

Al trabajar en el Paso 1: Analizar documentos con ai_parse_document(), 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 ayuda a mantener honestas las futuras modificaciones del código. Guarde 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. Mida la tasa de recuperación 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.

from pyspark.sql.functions import expr

# Volume path where PDFs are stored
docs_path = "/Volumes/<YOUR_CATALOG>/<YOUR_SCHEMA>/source_docs/"
# Read all files as binary
docs_df = spark.read.format("binaryFile").load(docs_path)
# Parse each document using ai_parse_document v2.0
parsed_df = docs_df.withColumn(
    "parsed_content",
    expr(f"""ai_parse_document(content, map(
        "version", "2.0",
        "imageOutputPath", "{docs_path}/parsed_images/"
    ))""")
)
# Drop binary content (too large to display)
parsed_df = parsed_df.drop("content")
# Save to Delta table
output_table = "<YOUR_CATALOG>.<YOUR_SCHEMA>.docs_parsed"
parsed_df.write.format("delta").mode("overwrite").saveAsTable(output_table)
print(f"✅ Parsed results saved to: {output_table}")

Paso 2: Limpiar, transformar y dividir en fragmentos

Al trabajar en el Paso 2: Limpiar, Transformar y Dividir en fragmentos, 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 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. Al trabajar en el Paso 2: Limpiar, Transformar y Dividir en fragmentos, 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.

Extracción rápida de texto plano

La extracción rápida de texto plano funciona mejor cuando se trata como una superficie medible. Capture un transcripte exitoso, un caso de fallo y la nota de reversión antes de ampliar el alcance. Registre los tiempos y el costo por token o consulta junto con los resultados funcionales. Tener visibilidad del costo desde temprano evita facturas inesperadas cuando el proceso pasa de la versión de demostración a entornos compartidos. 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.

from pyspark.sql import functions as F

# Convert VARIANT to JSON string, then extract text content
safe_json_col = F.coalesce(
    F.to_json(F.col("parsed_content")),
    F.col("parsed_content").cast("string")
)
plain_text_df = parsed_df.withColumn(
    "plain_text",
    extract_contents_udf()(safe_json_col)  # Custom UDF to join text elements
)

Fragmentar con LangChain

El chunking con LangChain 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. 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. Separe la política de chunking de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

from langchain_text_splitters import RecursiveCharacterTextSplitter
from pyspark.sql.types import StructType, StructField, StringType
import pandas as pd

CHUNK_SIZE = 2000
CHUNK_OVERLAP = 200
splitter = RecursiveCharacterTextSplitter(
    chunk_size=CHUNK_SIZE,
    chunk_overlap=CHUNK_OVERLAP,
    separators=["\n== page ==\n", "== page ==", "\n\n", "\n", " ", ""]
)
schema = StructType([
    StructField("path", StringType(), True),
    StructField("chunk", StringType(), True),
])
def split_rows(iterator):
    for pdf in iterator:
        out = []
        for _, row in pdf.iterrows():
            path, text = row["path"], row["plain_text"]
            if isinstance(text, str) and text.strip():
                for c in splitter.split_text(text):
                    if c and c.strip():
                        out.append((path, c))
        yield pd.DataFrame(out, columns=["path", "chunk"])
df_chunks = (
    plain_text_df.select("path", "plain_text")
    .mapInPandas(split_rows, schema=schema)
)
# Add unique IDs and save
df_chunks = df_chunks.withColumn("id", F.monotonically_increasing_id())
chunked_table = "<YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked"
df_chunks.write.format("delta") \
    .mode("overwrite") \
    .option("mergeSchema", "true") \
    .saveAsTable(chunked_table)

Paso 3: Construir la búsqueda vectorial

Paso 3: La búsqueda vectorial 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 mejoras posteriores. Separe la política de particionamiento de la política de recuperación; cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

Habilitar el Change Data Feed

Activar la función de Change Data Feed funciona mejor cuando se trata como una métrica cuantificable. Capture un registro de éxito 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 problema debe referirse a una sola responsabilidad y no a un proceso complejo e entrelazado. 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.

ALTER TABLE <YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);

Crear el índice de sincronización Delta

Crear el índice de sincronización Delta funciona mejor cuando se trata como una superficie medible. Capture una transcripción de referencia, un caso de fallo y la nota de reversión antes de ampliar el alcance. Considere esta etapa como un contrato entre las entradas y los resultados validados. 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.

from databricks.vector_search.client import VectorSearchClient

vsc = VectorSearchClient(disable_notice=True)
index_name = "<YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked_index"
vsc.create_delta_sync_index_and_wait(
    endpoint_name="<YOUR_VS_ENDPOINT>",
    index_name=index_name,
    source_table_name="<YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked",
    primary_key="id",
    embedding_source_column="chunk",
    embedding_model_endpoint_name="databricks-gte-large-en",
    pipeline_type="TRIGGERED",
)

Prueba de recuperación

La recuperación de pruebas funciona mejor cuando se trata como una métrica cuantificable. Capture un registro exitoso, 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 cuando el proceso pasa de la versión de demostración a entornos compartidos. Separe 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.

index = vsc.get_index(index_name=index_name)

results = index.similarity_search(
    query_text="How does the system prevent overheating?",
    columns=["path", "chunk"],
    num_results=5,
)
display(results)

Paso 4: Configurar Lakebase para la memoria de conversación

Paso 4: Configurar Lakebase para la memoria de conversaciones funciona mejor cuando se trata como una superficie medible. Capture un transcripte ejemplar, 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. 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.

Proveer un proyecto de escalado automático de Lakebase

Provision a Lakebase Autoscaling Project 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. 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 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. Provision a Lakebase Autoscaling Project 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 las completaciones parciales silenciosas.

Obtener detalles de conexión de forma programática

Para obtener los detalles de la conexión de forma programática, 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. Registre los tiempos de ejecución y el costo de tokens o consultas junto con los resultados funcionales. La visibilidad temprana del costo evita facturas inesperadas cuando el proceso pasa de entornos de demostración a entornos compartidos. Cite los pasajes que realmente sustentan la respuesta; sin citas, los operadores no pueden distinguir entre alucinaciones y lagunas en el indexado.

from databricks.sdk import WorkspaceClient

w = WorkspaceClient()
project_id = "<YOUR_PROJECT_NAME>"
# Get branch and endpoint
branches = list(w.postgres.list_branches(parent=f"projects/{project_id}"))
branch_name = branches[0].name
endpoints = list(w.postgres.list_endpoints(parent=branch_name))
ep = endpoints[0]
HOST = ep.status.hosts.host
ENDPOINT = ep.name
USERNAME = "<YOUR_DATABRICKS_EMAIL>"
print(f"Host: {HOST}")
print(f"Endpoint: {ENDPOINT}")

Probar la conexión

Para probar la conexió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 desde un punto de control conocido sin tener que adivinar el estado oculto. Guarde la configuración fuera del código de la aplicación. Los archivos de entorno, los almacenes de datos confidenciales y las banderas de funcionalidad deben encontrarse en un lugar donde los operadores puedan auditarlos sin necesidad de leer todo el sistema. Mencione las secciones que realmente sirvieron como base para la respuesta. Sin citaciones, los operadores no podrán distinguir entre una alucinación y una laguna en el indexado.

import psycopg2

cred = w.postgres.generate_database_credential(endpoint=ENDPOINT)
conn = psycopg2.connect(
    host=HOST,
    dbname="databricks_postgres",
    user=USERNAME,
    password=cred.token,
    port=5432,
    sslmode="require"
)
with conn.cursor() as cur:
    cur.execute("SELECT version()")
    print(cur.fetchone()[0])
conn.close()
print("✅ Connected to Lakebase!")

Crear tablas de puntos de control

Para crear tablas de puntos de control, 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.

from urllib.parse import quote
from langgraph.checkpoint.postgres import PostgresSaver

cred = w.postgres.generate_database_credential(endpoint=ENDPOINT)
DB_URI = (
    f"postgresql://{quote(USERNAME, safe='')}:{quote(cred.token, safe='')}"
    f"@{HOST}:5432/databricks_postgres"
    f"?sslmode=require"
)
with PostgresSaver.from_conn_string(DB_URI) as checkpointer:
    checkpointer.setup()
    print("✅ Checkpoint tables created!")

Para crear tablas de puntos de control, 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. Considere 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.

Paso 5: Construir el agente consciente del contexto

Al trabajar en el Paso 5: Construir el agente consciente del contexto, 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. 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 cuando se pasa de entornos de demostración a entornos compartidos. 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.

Agente interactivo (cuaderno de notas)

Al trabajar con el Interactive Agent (Notebook), 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 ayuda a mantener honestas las futuras modificaciones del código. Guarde 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. Mida la tasa 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.

from urllib.parse import quote
from langchain.agents import create_agent
from databricks_langchain import ChatDatabricks, VectorSearchRetrieverTool
from langgraph.checkpoint.postgres import PostgresSaver
from databricks.sdk import WorkspaceClient
import psycopg
from psycopg.rows import dict_row

def get_lakebase_checkpointer(host: str, endpoint: str, username: str):
    """Create a PostgresSaver backed by Lakebase Autoscaling."""
    w = WorkspaceClient()
    cred = w.postgres.generate_database_credential(endpoint=endpoint)
    db_uri = (
        f"postgresql://{quote(username, safe='')}:{quote(cred.token, safe='')}"
        f"@{host}:5432/databricks_postgres"
        f"?sslmode=require"
    )
    # IMPORTANT: Use psycopg.connect directly, not from_conn_string
    # from_conn_string returns a context manager, not a persistent instance
    conn = psycopg.connect(db_uri, autocommit=True, row_factory=dict_row)
    checkpointer = PostgresSaver(conn=conn)
    checkpointer.setup()
    return checkpointer

def build_agent(llm_endpoint: str, index_name: str, num_results: int = 3):
    model = ChatDatabricks(endpoint=llm_endpoint, max_tokens=500)
    vs_tool = VectorSearchRetrieverTool(
        name="knowledge_search",
        index_name=index_name,
        description="Search knowledge base for relevant information",
        num_results=num_results,
    )
    # Lakebase-backed checkpointer instead of InMemorySaver
    checkpointer = get_lakebase_checkpointer(HOST, ENDPOINT, USERNAME)
    system_prompt = """You are a Knowledge Assistant. Respond in a clear,
    professional tone. Use only verified information from the provided documents.
    If the answer cannot be found, clearly state that."""
    return create_agent(
        model=model,
        tools=[vs_tool],
        system_prompt=system_prompt,
        checkpointer=checkpointer,
    )

Probar conversaciones de múltiples turnos

Al trabajar en la prueba de conversación multietapa, 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 garantiza que los cambios posteriores en el código sean transparentes. 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 de mejoras posteriores. 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.

agent = build_agent("<YOUR_LLM_ENDPOINT>", "<YOUR_INDEX_NAME>", 3)

# STABLE thread_id - this is what enables context awareness
config = {"configurable": {"thread_id": "demo-session-001"}}
# Turn 1
r1 = agent.invoke(
    {"messages": [{"role": "user", "content": "What is the Orion system?"}]},
    config=config
)
print("Turn 1:", r1['messages'][-1].content)
# Turn 2 - agent should know "it" = Orion
r2 = agent.invoke(
    {"messages": [{"role": "user", "content": "How does it handle overheating?"}]},
    config=config
)
print("Turn 2:", r2['messages'][-1].content)

Al trabajar en la prueba de conversación multietapa, 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 garantiza que los cambios posteriores en el código sean transparentes. Considere esta etapa como un contrato entre las entradas y los resultados validados. Asigne nombres a los artefactos, defina las verificaciones de éxito y evite completaciones parciales silenciosas.

Paso 6: Código del agente de producción (agent.py)

Paso 6: El código del agente de producción (agent.py) 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. Registre los tiempos y el costo de tokens o consultas junto con los resultados funcionales. Tener visibilidad del costo desde el principio evita facturas inesperadas cuando el proceso pasa de entornos de demostración a entornos compartidos. 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.

# agent.py
import os
from uuid import uuid4
from typing import Any, Dict, List
from urllib.parse import quote

import yaml
import mlflow
import psycopg
from psycopg.rows import dict_row
from mlflow.pyfunc import ResponsesAgent
from mlflow.types.responses import ResponsesAgentRequest, ResponsesAgentResponse
from langchain.agents import create_agent
from databricks_langchain import ChatDatabricks, VectorSearchRetrieverTool
from langgraph.checkpoint.postgres import PostgresSaver
from databricks.sdk import WorkspaceClient

def _load_config(path: str = "agent-config.yaml") -> Dict[str, Any]:
    if not os.path.exists(path):
        raise FileNotFoundError(f"Config file not found at '{path}'")
    with open(path, "r", encoding="utf-8") as f:
        cfg = yaml.safe_load(f) or {}
    llm_endpoint = cfg.get("llm_endpoint_name")
    vs = cfg.get("vector_search", {}) or {}
    index_name = vs.get("index_name")
    num_results = int(vs.get("num_results", 3))
    lakebase = cfg.get("lakebase", {}) or {}
    return {
        "llm_endpoint_name": llm_endpoint,
        "vs_index_name": index_name,
        "vs_num_results": num_results,
        "lakebase_host": lakebase.get("host"),
        "lakebase_endpoint": lakebase.get("endpoint"),
        "lakebase_user": lakebase.get("user"),
    }

def get_lakebase_checkpointer(host, endpoint, user):
    w = WorkspaceClient()
    cred = w.postgres.generate_database_credential(endpoint=endpoint)
    db_uri = (
        f"postgresql://{quote(user, safe='')}:{quote(cred.token, safe='')}"
        f"@{host}:5432/databricks_postgres?sslmode=require"
    )
    conn = psycopg.connect(db_uri, autocommit=True, row_factory=dict_row)
    checkpointer = PostgresSaver(conn=conn)
    checkpointer.setup()
    return checkpointer

def build_agent(llm_endpoint, index_name, num_results,
                lakebase_host, lakebase_endpoint, lakebase_user):
    model = ChatDatabricks(endpoint=llm_endpoint, max_tokens=500)
    vs_tool = VectorSearchRetrieverTool(
        name="knowledge_search",
        index_name=index_name,
        description="Search knowledge base for relevant information",
        num_results=num_results,
    )
    checkpointer = get_lakebase_checkpointer(
        lakebase_host, lakebase_endpoint, lakebase_user
    )
    system_prompt = (
        "You are a Knowledge Assistant. Respond in a clear, professional tone. "
        "Use only verified information from the provided documents. "
        "If the answer cannot be found, clearly state that."
    )
    return create_agent(
        model=model, tools=[vs_tool],
        system_prompt=system_prompt, checkpointer=checkpointer,
    )

def _last_user_text(messages):
    user_msgs = [m for m in messages if m.get("role") == "user"]
    return str(user_msgs[-1].get("content", "")) if user_msgs else ""

class LangChainResponsesAgent(ResponsesAgent):
    def __init__(self):
        cfg = _load_config()
        self._agent = build_agent(
            cfg["llm_endpoint_name"], cfg["vs_index_name"],
            cfg["vs_num_results"], cfg["lakebase_host"],
            cfg["lakebase_endpoint"], cfg["lakebase_user"],
        )
    def predict(self, request: ResponsesAgentRequest) -> ResponsesAgentResponse:
        msgs = [m.model_dump() for m in request.input]
        custom_inputs = dict(request.custom_inputs or {})
        thread_id = custom_inputs.get("thread_id", f"session-{uuid4()}")
        result = self._agent.invoke(
            {"messages": msgs},
            config={"configurable": {"thread_id": thread_id}},
        )
        try:
            text = result["messages"][-1].content
        except Exception:
            text = str(result)
        return ResponsesAgentResponse(
            output=[self.create_text_output_item(text, str(uuid4()))],
            custom_outputs={"thread_id": thread_id},
        )

AGENT = LangChainResponsesAgent()
mlflow.models.set_model(AGENT)

Configuración (agent-config.yaml)

La configuración (agent-config.yaml) funciona mejor cuando se trata como un elemento medible. Capture una transcripción de referencia, 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. Separe la política de particionamiento de la política de recuperación. Cambiar una no debe obligar a reescribir la otra cuando cambian las métricas de calidad.

llm_endpoint_name: <YOUR_LLM_ENDPOINT>
vector_search:
  index_name: <YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked_index
  num_results: 3
lakebase:
  host: <YOUR_LAKEBASE_HOST>
  endpoint: projects/<YOUR_PROJECT>/branches/production/endpoints/primary
  user: <YOUR_DATABRICKS_EMAIL>

Paso 7: Registrar, registrar y desplegar

Paso 7: Registrar, registrar y desplegar funciona mejor cuando se trata como un elemento medible. Capture una transcripción de referencia, un caso de fallo y la nota de reversión antes de ampliar el alcance.

Registrar en MLflow

import mlflow
from importlib.metadata import version as get_version
from mlflow.models.resources import DatabricksVectorSearchIndex, DatabricksServingEndpoint

resources = [
    DatabricksVectorSearchIndex(index_name="<YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked_index"),
    DatabricksServingEndpoint(endpoint_name="<YOUR_LLM_ENDPOINT>"),
]
with mlflow.start_run():
    mlflow.set_tags({
        "model_type": "retrieval_agent",
        "framework": "langchain",
        "memory": "lakebase_autoscaling",
    })
    logged_agent_info = mlflow.pyfunc.log_model(
        name="knowledge_assistant",
        python_model="agent.py",
        code_paths=["agent-config.yaml"],
        input_example={"input": [{"role": "user", "content": "What is Orion?"}]},
        pip_requirements=[
            f"databricks-vectorsearch=={get_version('databricks-vectorsearch')}",
            f"databricks-langchain=={get_version('databricks-langchain')}",
            f"langchain=={get_version('langchain')}",
            f"mlflow=={get_version('mlflow')}",
            "langgraph-checkpoint-postgres",
            "psycopg[binary]",
            "databricks-sdk>=0.89.0",
        ],
        resources=resources,
    )
    model_uri = logged_agent_info.model_uri

Registrar en Unity Catalog

mlflow.set_registry_uri("databricks-uc")
UC_MODEL_NAME = "<YOUR_CATALOG>.<YOUR_SCHEMA>.knowledge_assistant"

uc_info = mlflow.register_model(model_uri=model_uri, name=UC_MODEL_NAME)
print(f"✅ Registered: {UC_MODEL_NAME} v{uc_info.version}")

Desplegar

from databricks import agents

deployment = agents.deploy(
    model_name=UC_MODEL_NAME,
    model_version=uc_info.version,
    scale_to_zero_enabled=True,
)
print(f"✅ Endpoint: {deployment.query_endpoint}")

El resultado

ws = WorkspaceClient()
client = ws.serving_endpoints.get_open_ai_client()

session = "user-session-042"
# Turn 1
r1 = client.responses.create(
    model="knowledge_assistant",
    input=[{"role": "user", "content": "What is the Orion motion controller?"}],
    extra_body={"custom_inputs": {"thread_id": session}}
)
# Turn 2 - "it" resolves correctly to Orion
r2 = client.responses.create(
    model="knowledge_assistant",
    input=[{"role": "user", "content": "How does it prevent overheating?"}],
    extra_body={"custom_inputs": {"thread_id": session}}
)

Verificación rápida para validar la persistencia de la memoria:

import json
from uuid import uuid4

thread_id = f"memory-test-{uuid4()}"

# --- 1. Use custom input instead of input_example ---
custom_input = {
    "input": [{"role": "user", "content": "What are the main components of Orion?"}],
    "custom_inputs": {"thread_id": thread_id},
}

print("=" * 50)
print("TEST 1: Custom input (no input_example needed)")
print("=" * 50)

# Use output_path to capture results (without it, mlflow.models.predict returns None)
mlflow.models.predict(
    model_uri=model_uri,
    input_data=custom_input,
    env_manager="uv",
    output_path="/tmp/result_1.json",
)

with open("/tmp/result_1.json", "r") as f:
    result_1 = json.load(f)

thread_id_1 = result_1["custom_outputs"]["thread_id"]
response_1 = result_1["output"][0]["content"][0]["text"]
print(f"Thread ID: {thread_id_1}")
print(f"Response: {response_1[:300]}...")

# --- 2. Test memory persistence with same thread_id ---
follow_up_input = {
    "input": [
        {
            "role": "user",
            "content": "Can you elaborate more on the first component you mentioned?",
        }
    ],
    "custom_inputs": {"thread_id": thread_id},  # Reuse same thread
}

print("\n" + "=" * 50)
print("TEST 2: Follow-up on same thread (memory test)")
print("=" * 50)

mlflow.models.predict(
    model_uri=model_uri,
    input_data=follow_up_input,
    env_manager="uv",
    output_path="/tmp/result_2.json",
)

with open("/tmp/result_2.json", "r") as f:
    result_2 = json.load(f)

thread_id_2 = result_2["custom_outputs"]["thread_id"]
response_2 = result_2["output"][0]["content"][0]["text"]
print(f"Thread ID: {thread_id_2}")
print(f"Response: {response_2[:300]}...")

# --- 3. Verify memory with actual conditions ---
print("\n" + "=" * 50)
print("MEMORY CHECK")
print("=" * 50)

# Check 1: Thread IDs match
if thread_id_1 == thread_id_2:
    print(f"✅ Thread ID match: {thread_id_1}")
else:
    print(f"❌ Thread ID mismatch! Call 1: {thread_id_1}, Call 2: {thread_id_2}")

# Check 2: Follow-up response references context from the first response
follow_up_lower = response_2.lower()
if len(response_2) > 50 and any(
    keyword in follow_up_lower
    for keyword in [
        "motion",
        "vision",
        "cognition",
        "communication",
        "subsystem",
        "component",
    ]
):
    print(
        "✅ Follow-up response references components from the first answer — memory is intact!"
    )
else:
    print(
        "⚠️ Follow-up response may not reference the first answer. Manual review recommended."
    )
    print(f"   Follow-up preview: {response_2[:200]}")

print(
    "\n✅ Lakebase Postgres checkpointing is working correctly!"
    if thread_id_1 == thread_id_2
    else "\n❌ Memory persistence test FAILED."
)

Capturas de pantalla de las tablas de puntos de control de Lakebase Postgres:

Dificultades que se pueden encontrar en el proceso

Qué cambió respecto a un agente RAG estándar

Conclusión

Lista de verificación operativa