Startseite / Artikel / Ihr RAG-Agent vergisst nach einer Nachricht alles – So haben Sie das mit Databricks behoben…

Ihr RAG-Agent vergisst nach einer Nachricht alles – So haben Sie das mit Databricks behoben…

Schritt-für-Schritt-Anleitung: Ihr RAG-Agent vergisst nach einer Nachricht alles – So haben Sie das mit Databricks behoben…: Verträge, Überprüfungen sowie Code-Blöcke für Teams.

3410 Wörter

Dieser Leitfaden zeigt Schritt für Schritt den Weg von Rohstoffen bis zu einem funktionsfähigen System für folgendes Problem: Ihr RAG-Agent vergisst nach einer Nachricht alles – hier zeigt sich, wie ich das mit Databricks Lakebase behoben habe. Der Schwerpunkt liegt auf umsetzbaren Schritten, expliziten Überprüfungen sowie Code, den Sie ohne Rückschluss auf die Absicht direkt in ein Repository einfügen können. Zur Übersicht sollten Sie vor der Änderung des Codes die Eingaben, den Verantwortlichen für den Schritt sowie die Abbruchkriterien definieren. Die Operator sollten in der Lage sein, den Schritt von einem bekannten Checkpoint aus erneut auszuführen, ohne auf versteckte Zustände schließen zu müssen. Betrachten Sie diese Phase als Vertrag zwischen den Eingaben und den validierten Ausgaben. Benennen Sie die Artefakte, definieren Sie Erfolgskontrollen und lehnen Sie stille, unvollständige Abschlüsse ab.

Die Architektur

Beim Arbeiten an „The Architecture“ sollten Sie zunächst den Vertrag aufschreiben: erforderliche Eingaben, Erfolgsignal sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Notieren Sie die Laufzeiten sowie die Kosten für Token oder Abfragen neben den funktionalen Ergebnissen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn der Einsatzbereich von einer Demo in gemeinsame Umgebungen wechselt. Messen Sie die Trefferquote anhand einer festgelegten Fragestellung, bevor Sie die Anfragen anpassen. Ein häufiges Wechseln der Anfragemuster behebt in der Regel nicht ein schwaches Suchverhalten.

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)

Schritt 1: Dokumente mit ai_parse_document() parsen

Beim Bearbeiten von Schritt 1: Dokumente mit ai_parse_document() parsen, sollten Sie zunächst den Vertrag aufschreiben: erforderliche Eingaben, Erfolgsignal sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Bewahren Sie die Konfiguration außerhalb des Anwendungscode auf. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort gesammelt sein, den Betreiber ohne das Durchlesen des gesamten Systems überprüfen können. Messen Sie die Trefferquote anhand einer festgelegten Fragestellung, bevor Sie die Anfragen anpassen. Eine häufige Änderung der Anfragen behebt selten ein schwaches Extraktionsverhalten.

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

Schritt 2: Reinigen, Transformieren und in Blöcke aufteilen

Beim Bearbeiten von Schritt 2: Reinigen, Transformieren und in Blöcke aufteilen, sollten Sie zunächst den Vertrag festhalten: erforderliche Eingaben, Erfolgsignal sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Dokumentieren Sie sowohl den erfolgreichen Ablauf als auch den Notfallweg gemeinsam. Wiederholte Versuche, menschliche Überprüfungen sowie die Handhabung von Fehlern gehören zum Produkt selbst und nicht zu späteren Optimierungen. Messen Sie die Trefferquote anhand eines festgelegten Fragebogens, bevor Sie die Anfragen anpassen – ein häufiges Wechseln der Anfragen behebt selten eine schwache Suchfunktion. Behandeln Sie diese Phase als Vertrag zwischen Eingaben und validierten Ausgaben: Nennen Sie die Erzeugnisse, definieren Sie Erfolgskontrollen und lehnen Sie stille, teilweise abgeschlossene Ergebnisse ab.

Schnelle Extraktion von reinem Text

Fast Plain Text Extraction funktioniert am besten, wenn sie als messbare Ebene betrachtet wird. Erfassen Sie ein gelungenes Beispiel, einen Fehlfall sowie die Notizen zur Rücksetzung, bevor Sie den Umfang erweitern. Erfassen Sie außerdem die Laufzeiten sowie die Kosten pro Token oder Abfrage zusammen mit den funktionalen Ergebnissen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn sich der Einsatzbereich von einer Demo auf gemeinsame Umgebungen verschiebt. Trennen Sie die Strategie zur Aufteilung in Blöcke von der Strategie zur Abrufung. Ein Änderungsbedarf bei einer dieser Strategien sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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
)

Blöcke mit LangChain erstellen

Chunking mit LangChain funktioniert am besten, wenn es als messbare Struktur betrachtet wird. Erfassen Sie ein optimales Beispiel, einen Fehlerfall sowie eine Notiz zur Rücksetzung, bevor Sie den Umfang erweitern. Bewahren Sie die Konfiguration außerhalb des Anwendungscode auf. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort gesammelt sein, den Betreuer ohne Durchsicht des gesamten Systems prüfen können. Trennen Sie die Chunking-Strategie von der Abrufstrategie. Ein Änderung in einer sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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)

Schritt 3: Vector Search erstellen

Schritt 3: Die Implementierung einer Vektor-Suche funktioniert am besten, wenn sie als messbare Größe betrachtet wird. Erfassen Sie vor der Erweiterung des Umfangs ein erfolgreiches Beispiel, einen Fehlerfall sowie die Notizen zur Rücksetzung. Dokumentieren Sie gleichzeitig den erfolgreichen Ablauf sowie den Wiederherstellungsprozess. Wiederholversuche, menschliche Überprüfungen und die Handhabung von Fehlern gehören zum Produkt selbst, nicht zu späteren Optimierungen. Trennen Sie die Strategie zur Aufteilung in Blöcke von der Strategie zur Suche. Ein Änderung einer dieser Strategien sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

Change Data Feed aktivieren

Die Aktivierung von Change Data Feed funktioniert am besten, wenn sie als messbarer Ansatz betrachtet wird. Erfassen Sie ein „goldenes“ Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung, bevor Sie den Umfang erweitern. Ziehen Sie kleine, testbare Einheiten vor großen, komplexen Skripten vor. Wenn ein Schritt fehlschlägt, sollte der Fehler auf eine einzige Verantwortung verweisen und nicht auf einen verworrenen Ablauf. Trennen Sie die Chunking-Strategie von der Abrufstrategie – eine Änderung sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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

Erstellen Sie den Delta Sync Index

Die Erstellung des Delta Sync Index funktioniert am besten, wenn sie als messbare Größe betrachtet wird. Erfassen Sie ein „goldenes“ Transkript, einen Fehlerfall sowie die Rollback-Anmerkung, bevor Sie den Umfang erweitern. Betrachten Sie diese Phase als Vertrag zwischen den Eingaben und den validierten Ausgaben. Benennen Sie die Artefakte, definieren Sie Erfolgskontrollen und lehnen Sie stille, teilweise abgeschlossene Arbeiten ab. Trennen Sie die Chunking-Strategie von der Abrufstrategie. Ein Änderung in einer sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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

Test des Abrufs

Die Testabfrage funktioniert am besten, wenn sie als messbarer Ansatz betrachtet wird. Erfassen Sie ein erfolgreiches Beispiel, einen Fehlerfall sowie die Notizen zur Rücksetzung, bevor Sie den Umfang erweitern. Erfassen Sie außerdem die Laufzeiten sowie die Kosten für Tokens oder Abfragen neben den funktionalen Ergebnissen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn sich der Einsatzbereich von einer Demo-Umgebung auf gemeinsam genutzte Umgebungen verschiebt. Trennen Sie die Strategie zur Aufteilung in Blöcke von der Abfragestrategie – ein Änderungsbedarf bei einer Seite sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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)

Schritt 4: Lakebase für das Gesprächsgedächtnis einrichten

Schritt 4: Einrichtung von Lakebase für das Gesprächsmerkmals-Management funktioniert am besten, wenn es als messbare Struktur betrachtet wird. Erfassen Sie vor der Erweiterung des Umfangs ein gelungenes Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung. Bewahren Sie die Konfiguration außerhalb des Anwendungscode auf. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort gesammelt sein, den Betreiber ohne das Durchlesen des gesamten Systems prüfen können. Trennen Sie die Chunking-Strategie von der Abrufstrategie. Ein Änderung in einer sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

Einrichten eines Lakebase-Autoscaling-Projekts

Ein Lakebase Autoscaling Project funktioniert am besten, wenn es als messbarer Bereich betrachtet wird. Erfassen Sie vor der Erweiterung des Umfangs ein „goldenes“ Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung. Dokumentieren Sie gemeinsam den erfolgreichen Ablauf sowie den Wiederherstellungsprozess. Wiederholversuche, menschliche Überprüfungen und die Handhabung von Fehlern gehören zum Produkt selbst, nicht zu späteren Optimierungen. Trennen Sie die Chunking-Strategie von der Abrufstrategie. Eine Änderung sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern. Ein Lakebase Autoscaling Project funktioniert am besten, wenn es als messbarer Bereich betrachtet wird. Erfassen Sie vor der Erweiterung des Umfangs ein „goldenes“ Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung. Betrachten Sie diese Phase als Vertrag zwischen den Eingaben und den validierten Ausgaben. Benennen Sie die Artefakte, definieren Sie Erfolgskontrollen und lehnen Sie stille, unvollständige Abschlüsse ab.

Verbindungsdaten programmgesteuert abrufen

Für die programmatische Abfrage von Verbindungsdaten sollten vor der Codeänderung Eingabedaten, der Eigentümer des Schritts sowie die Abbruchkriterien definiert werden. Die Operator sollten in der Lage sein, den Schritt von einem bekannten Checkpoint aus erneut auszuführen, ohne auf versteckten Zuständen schließen zu müssen. Erhalten Sie neben den funktionalen Ergebnissen auch Aufzeichnungen der Laufzeiten sowie der Kosten für Token oder Abfragen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn der Ablauf von Demonstrumgebungen in gemeinsam genutzte Umgebungen wechselt. Zitieren Sie die Passagen, die tatsächlich die Antwort untermauern. Ohne Zitate können die Operator nicht zwischen Halluzinationen und Lücken in der Indizierung unterscheiden.

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

Verbindung testen

Zum Testen der Verbindung sollten die Eingaben, der Verantwortliche für den Schritt sowie die Abbruchkriterien vor dem Ändern des Codes definiert werden. Die Operator sollten in der Lage sein, den Schritt von einem bekannten Checkpoint aus erneut auszuführen, ohne auf versteckte Zustände schließen zu müssen. Die Konfiguration sollte außerhalb des Anwendungscode gespeichert werden. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort zusammengefasst sein, den die Operator überprüfen können, ohne den gesamten Ablauf durchzulesen. Zitieren Sie die Passagen, die tatsächlich die Antwort untermauern. Ohne Zitate können die Operator nicht zwischen Halluzinationen und Lücken in der Indizierung unterscheiden.

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

Erlernen Sie, Checkpoint-Tabellen zu erstellen

Für die Erstellung von Checkpoint-Tabellen sollten vor der Codeänderung die Eingaben, der Verantwortliche für den Schritt sowie die Abbruchkriterien definiert werden. Die Bediener sollten in der Lage sein, den Schritt von einem bekannten Checkpoint aus erneut auszuführen, ohne auf versteckte Zustände schließen zu müssen. Dokumentieren Sie gemeinsam den erfolgreichen Ablauf sowie den Notfallweg. Wiederholungsversuche, menschliche Überprüfungen und die Handhabung von Fehlern gehören zum Produkt selbst und nicht zu späteren Optimierungen. Zitieren Sie die Passagen, die tatsächlich die Antwort begründen. Ohne Zitate können die Bediener nicht zwischen Halluzinationen und Lücken in der Indizierung unterscheiden.

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

Für die Erstellung von Checkpoint-Tabellen sollten vor der Codeänderung die Eingaben, der Verantwortliche für den Schritt sowie die Abbruchkriterien definiert werden. Die Bediener sollten in der Lage sein, den Schritt von einem bekannten Checkpoint aus erneut auszuführen, ohne auf versteckte Zustände schließen zu müssen. Betrachten Sie diese Phase als Vertrag zwischen den Eingaben und den validierten Ausgaben. Benennen Sie die Ergebnisdokumente, definieren Sie Erfolgskontrollen und lehnen Sie stille, unvollständige Abschlüsse ab.

Schritt 5: Erstellen Sie den kontextbewussten Agenten

Beim Arbeiten am Schritt 5: Erstellen Sie den kontextbewussten Agenten, notieren Sie zunächst die Anforderungen: erforderliche Eingaben, Erfolgsindikator sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Notieren Sie außerdem die Laufzeiten sowie die Kosten für Tokens oder Abfragen neben den funktionalen Ergebnissen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn der Einsatzbereich von einer Demo in gemeinsame Umgebungen wechselt. Messen Sie die Genauigkeit anhand einer festgelegten Fragestellung, bevor Sie die Anfragen anpassen. Das häufige Wechseln der Anfragen behebt selten ein schwaches Suchsystem.

Interaktiver Agent (Notebook)

Beim Arbeiten am Interactive Agent (Notebook) sollten Sie zunächst den Vertrag aufschreiben: erforderliche Eingaben, Erfolgsignal sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Bewahren Sie die Konfiguration außerhalb des Anwendungscode auf. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort gesammelt sein, den Betreiber ohne das Durchlesen des gesamten Graphen überprüfen können. Messen Sie die Recall-Rate anhand eines festgelegten Fragekatalogs, bevor Sie die Prompts anpassen. Eine häufige Änderung der Prompts behebt selten ein schwaches Retrieval-System.

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,
    )

Teste mehrstufige Konversationen

Beim Arbeiten an Test-Mehrfachgesprächen sollten Sie zunächst den Vertrag aufschreiben: erforderliche Eingaben, Erfolgszeichen sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Dokumentieren Sie gleichzeitig den erfolgreichen Ablauf sowie den Notfallweg. Wiederholte Versuche, menschliche Überprüfungen und die Handhabung von Fehlern gehören zum Produkt selbst, nicht zu späteren Optimierungen. Messen Sie die Erinnerungsfähigkeit anhand eines festgelegten Fragebogens, bevor Sie die Anfragen anpassen. Eine häufige Änderung der Anfragen behebt selten ein schwaches Suchverhalten.

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)

Beim Arbeiten an Test-Mehrfachgesprächen sollten Sie zunächst den Vertrag aufschreiben: erforderliche Eingaben, Erfolgszeichen sowie das Vorgehen bei teilweisen Fehlern. Diese Checkliste sorgt dafür, dass spätere Codeänderungen transparent bleiben. Betrachten Sie diese Phase als Vertrag zwischen den Eingaben und den validierten Ausgaben. Benennen Sie die Ergebnisse, definieren Sie Erfolgskontrollen und lehnen Sie stille, teilweise abgeschlossene Abläufe ab.

Schritt 6: Code des Produktions-Agenten (agent.py)

Schritt 6: Der Code des Produktions-Agenten (agent.py) funktioniert am besten, wenn er als messbarer Ansatz betrachtet wird. Erfassen Sie vor der Erweiterung des Umfangs ein optimales Beispiel, einen Fehlerfall sowie eine Notiz zur Rücksetzung. Erfassen Sie außerdem die Laufzeiten sowie die Kosten für Tokens oder Abfragen neben den funktionalen Ergebnissen. Eine frühzeitige Sichtbarkeit der Kosten verhindert überraschende Rechnungen, wenn der Einsatzbereich von einer Demo-Umgebung in gemeinsam genutzte Umgebungen wechselt. Trennen Sie die Strategie zur Aufteilung in Blöcke von der Strategie zur Abrufung. Ein Änderungsbedarf bei einer dieser Strategien sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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

Konfiguration (agent-config.yaml)

Die Konfiguration (agent-config.yaml) funktioniert am besten, wenn sie als messbarer Bereich betrachtet wird. Erfassen Sie ein „goldenes“ Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung, bevor Sie den Umfang erweitern. Bewahren Sie die Konfiguration außerhalb des Anwendungscode auf. Umgebungsdateien, Geheimdatenspeicher und Feature-Flags sollten an einem Ort gesammelt sein, den Betreuer ohne das Durchlesen des gesamten Systems überprüfen können. Trennen Sie die Chunking-Strategie von der Abrufstrategie. Ein Änderung in einer sollte nicht dazu führen, dass die andere neu geschrieben werden muss, wenn sich die Qualitätsmetriken ändern.

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>

Schritt 7: Protokollieren, registrieren und bereitstellen

Schritt 7: Protokollieren, registrieren und bereitstellen funktioniert am besten, wenn er als messbarer Bereich betrachtet wird. Erfassen Sie ein „goldenes“ Transkript, einen Fehlerfall sowie eine Notiz zur Rücksetzung, bevor Sie den Umfang erweitern.

Protokollieren in 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

Registrieren im 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}")

Bereitstellen

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

Das Ergebnis

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

Schnelle Überprüfung zur Validierung der Speicherdauerhaftigkeit:

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

Screenshots der CheckPoint-Tabellen von Lakebase Postgres:

Häufige Probleme auf dem Weg

Was sich im Vergleich zu einem Standard-RAG-Agenten geändert hat

Zusammenfassung

Operative Checkliste