Twój agent RAG zapomina wszystko po jednej wiadomości – oto jak naprawiłeś to za pomocą Databricks…
Krok po kroku instrukcja obsługi problemu, gdy Twój agent RAG zapomina wszystko po jednej wiadomości – oto jak naprawiłeś to za pomocą Databricks…: umowy, sprawdzenia oraz gotowe miejsca na kod dla zespołów.
To przewodnik pokazuje, jak odbudować proces od surowców do działającego systemu w przypadku sytuacji, gdy Twój agent RAG zapomina wszystko po jednej wiadomości – oto jak naprawiłem to za pomocą Databricks Lakebase. Skupiamy się na krokach operacyjnych, wyraźnych sprawdzeniach oraz kodzie, który można bez problemu umieścić w repozytorium, bez konieczności domyślania się intencji. Aby uzyskać ogólny obraz, zdefiniuj wprowadzenia, osobę odpowiedzialną za dany krok oraz kryteria zakończenia przed modyfikacją kodu. Operatorzy powinni móc ponownie uruchomić ten krok na podstawie znanego punktu kontrolnego, bez konieczności zgadywania ukrytego stanu. Traktuj tę fazę jako umowę pomiędzy wprowadzeniami a zweryfikowanymi wynikami. Nadaj nazwy plikom, zdefiniuj kryteria sukcesu i odrzucaj ciche, częściowe ukończenie zadania.
Architektura
Gdy pracujesz nad „Architekturą”, najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Obok wyników funkcjonalnych zapisz czas wykonywania oraz koszt tokena lub zapytania. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy ścieżka przechodzi z wersji demonstracyjnej do środowisk współdzielonych. Zmierz stopień odzyskiwania informacji na ustalonej grupie pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko poprawiają słabą skuteczność wyszukiwania.
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)
Krok 1: Analiza dokumentów za pomocą ai_parse_document()
Gdy przechodzisz przez Krok 1: Analiza dokumentów za pomocą ai_parse_document(), najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie.
Przechowuj konfigurację poza kodem aplikacji. Pliki środowiskowe, magazyny tajnych danych oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury.
Zmierz stopień odzyskiwania informacji na ustalonej grupie pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko naprawiają słabe możliwości wyszukiwania.
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}")
Krok 2: Czyszczenie, transformacja i dzielenie na fragmenty
Gdy przechodzisz przez Krok 2: Oczyszczanie, transformacja i dzielenie na fragmenty, najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Zdokumentuj zarówno ścieżkę prawidłowego działania, jak i ścieżkę naprawczą. Próby ponowne, kontrola przez ludzi oraz obsługa wiadomości nieodebranych stanowią część produktu, a nie elementy dopinane później. Zmierz stopień odzyskiwania informacji na ustalonej grupie pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko naprawiają słabe możliwości wyszukiwania. Gdy przechodzisz przez Krok 2: Oczyszczanie, transformacja i dzielenie na fragmenty, najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Traktuj tę fazę jako umowę pomiędzy danymi wejściowymi a zweryfikowanymi wynikami. Nadaj nazwy artefaktom, zdefiniuj kryteria sukcesu i odrzucaj ciche, częściowe ukończenie zadań.
Szybkie wydobywanie prostego tekstu
Szybka ekstrakcja prostego tekstu działa najlepiej, gdy traktuje się ją jako mierzalną powierzchnię. Zapisz jeden idealny wynik, jeden przypadek awarii oraz notatkę o cofnięciu działań przed rozszerzaniem zakresu. Zarejestruj czasy wykonywania operacji oraz koszt tokenów lub zapytań obok wyników funkcjonalnych. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy przechodzi się z środowiska demonstracyjnego do współdzielonych środowisk. Oddziel zasadę dzielenia na fragmenty od zasady wyszukiwania. Zmiana jednej nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
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
)
Dzielenie na fragmenty za pomocą LangChain
Chunking z LangChain działa najlepiej, gdy traktuje się go jako mierzalną powierzchnię. Zapisz jeden idealny przykład transkrypcji, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres pracy. Przechowuj konfigurację poza kodem aplikacji. Pliki środowiskowe, składysek tajnych danych oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury. Oddziel zasadę dzielenia na chunki od zasady wyszukiwania. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
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)
Krok 3: Budowa wyszukiwania wektorowego
Krok 3: Budowa systemu wyszukiwania wektorowego działa najlepiej, gdy traktuje się go jako mierzalną powierzchnię. Zanim rozszerzysz zakres, zapisz jeden idealny przypadek działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia działań. Zdokumentuj zarówno prawidłowy przebieg operacji, jak i ścieżkę przywracania do normalnego stanu. Próby ponownych działań, kontrola przez ludzi oraz obsługa wiadomości nieodebranych stanowią część produktu, a nie elementy dodawane później. Oddziel zasadę dzielenia na fragmenty od zasady wyszukiwania. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
Włącz Change Data Feed
Aktywacja funkcji Change Data Feed działa najlepiej, gdy traktuje się ją jako mierzalną zmienną. Zapisz jeden idealny przykład działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres działania. Wolno preferować małe, testowalne jednostki zamiast rozbudowanych skryptów. Gdy jakiś krok zawiedzie, awaria powinna wskazywać na konkretną odpowiedzialność, a nie na skomplikowany łańcuch operacji. Rozdziel politykę dzielenia na fragmenty od polityki pobierania danych. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
ALTER TABLE <YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
Stwórz indeks synchronizacji Delta
Tworzenie indeksu Delta Sync działa najlepiej, gdy traktuje się go jako mierzalną powierzchnię. Zapisz jeden idealny przepis, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres pracy. Traktuj tę fazę jako umowę pomiędzy danymi wejściowymi a zweryfikowanymi wynikami. Nadaj nazwy poszczególnym elementom, zdefiniuj kryteria sukcesu i odrzucaj ciche, częściowe ukończenie zadań. Rozdziel politykę dzielenia na fragmenty od polityki wyszukiwania. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
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 wyszukiwania
Test pobierania danych działa najlepiej, gdy traktuje się go jako mierzalną powierzchnię do analizy. Zapisz jeden idealny przekaz, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia działań, zanim rozszerzysz zakres badania. Zarejestruj czasy wykonywania operacji oraz koszt tokenów lub zapytań obok wyników funkcjonalnych. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy przechodzi się z środowiska demonstracyjnego do współdzielonych środowisk. Oddziel zasady dzielenia danych na fragmenty od zasad pobierania danych – zmiana jednych nie powinna zmuszać do przepisywania drugich, gdy zmieniają się metryki jakości.
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)
Krok 4: Konfiguracja Lakebase do przechowywania pamięci rozmów
Krok 4: Konfiguracja Lakebase dla pamięci rozmów działa najlepiej, gdy jest traktowana jako mierzalna powierzchnia. Zapisz jeden idealny przepis działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres. Przechowuj konfigurację poza kodem aplikacji. Pliki środowiskowe, skrytki z danymi poufnymi oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury. Oddziel zasadę dzielenia na fragmenty od zasady wyszukiwania. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
Utwórz projekt z automatycznym skalowaniem w Lakebase
Provision a Lakebase Autoscaling Project funkcjonuje najlepiej, gdy jest traktowany jako mierzalna powierzchnia do analizy. Zapisz jeden idealny przepis działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian przed rozszerzaniem zakresu. Zdokumentuj zarówno prawidłowy przebieg działania, jak i ścieżkę przywracania do stanu poprzedniego. Próby ponownych działań, kontrolne punkty ludzkie oraz obsługa wiadomości błędowych stanowią część produktu, a nie elementy dodawane później. Oddziel zasadę dzielenia na fragmenty od zasady pobierania danych. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej w momencie zmian wskaźników jakości. Provision a Lakebase Autoscaling Project funkcjonuje najlepiej, gdy jest traktowany jako mierzalna powierzchnia do analizy. Zapisz jeden idealny przepis działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian przed rozszerzaniem zakresu. Traktuj ten etap jako umowę pomiędzy danymi wejściowymi a zweryfikowanymi wynikami. Nadaj nazwy poszczególnym elementom, zdefiniuj kryteria sukcesu i odrzucaj ciche, częściowe ukończenie zadań.
Pobieranie informacji o połączeniu w sposób programowy
Aby uzyskać informacje o połączeniu w sposób programowy, należy zdefiniować dane wejściowe, osobę odpowiedzialną za dany krok oraz kryteria zakończenia przed modyfikacją kodu. Operatorzy powinni móc ponownie uruchomić ten krok od znanego punktu kontrolnego, bez konieczności zgadywania ukrytego stanu. Należy rejestrować czasy wykonywania oraz koszt tokenów lub zapytań obok wyników funkcjonalnych. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy ścieżka przechodzi z środowiska demonstracyjnego do współdzielonych środowisk. Należy podawać fragmenty tekstu, które faktycznie stanowiły podstawę odpowiedzi. Bez tych odniesień operatorzy nie mogą odróżnić halucynacji od luki w indeksowaniu.
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}")
Testowanie połączenia
Aby przetestować połączenie, należy zdefiniować dane wejściowe, osobę odpowiedzialną za dany krok oraz kryteria zakończenia przed modyfikacją kodu. Operatorzy powinni móc ponownie uruchomić dany krok na podstawie znanego punktu kontrolnego, bez konieczności zgadywania ukrytego stanu. Konfigurację należy przechowywać poza kodem aplikacji. Pliki środowiskowe, magazyny tajnych danych oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury. Należy podawać konkretne fragmenty tekstu, na których opiera się odpowiedź. Bez tych odniesień operatorzy nie będą w stanie odróżnić halucynacji od braku danych w indeksie.
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!")
Tworzenie tabel punktów kontrolnych
Aby utworzyć tabele punktów kontrolnych, należy zdefiniować dane wejściowe, osobę odpowiedzialną za dany krok oraz kryteria zakończenia przed modyfikacją kodu. Operatorzy powinni móc ponownie uruchomić ten krok na podstawie znanego punktu kontrolnego, bez konieczności zgadywania ukrytego stanu. Należy udokumentować zarówno prawidłowy przebieg procesu, jak i ścieżkę naprawczą. Próby ponownych działań, kontrola przez ludzi oraz obsługa wiadomości błędnych stanowią część produktu, a nie elementy dodawane później. Należy podać fragmenty tekstu, które faktycznie stanowią podstawę odpowiedzi. Bez tych odniesień operatorzy nie będą w stanie odróżnić halucynacji od luki w indeksowaniu.
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!")
Aby utworzyć tabele punktów kontrolnych, należy zdefiniować dane wejściowe, osobę odpowiedzialną za dany krok oraz kryteria zakończenia przed modyfikacją kodu. Operatorzy powinni móc ponownie uruchomić ten krok na podstawie znanego punktu kontrolnego, bez konieczności zgadywania ukrytego stanu. Traktuj tę fazę jako umowę pomiędzy danymi wejściowymi a zweryfikowanymi wynikami. Należy nadać nazwy poszczególnym elementom, zdefiniować kryteria sukcesu oraz odrzucić ciche, częściowe ukończenie zadania.
Krok 5: Stworzenie agenta świadomego kontekstu
Podczas pracy nad Krokiem 5: Stworzenie agenta świadomego kontekstu, najpierw zapisz specyfikację: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Zapisz czas wykonywania oraz koszt tokenów lub zapytań obok wyników funkcjonalnych. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy przechodzi się z środowiska demonstracyjnego do współdzielonych środowisk. Zmierz stopień przywoływania informacji na ustalonej serii pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko poprawiają słabe możliwości wyszukiwania.
Agent interaktywny (Notatnik)
Gdy pracujesz z Interactive Agent (Notebook), najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Przechowuj konfigurację poza kodem aplikacji. Pliki środowiskowe, magazyny tajnych danych oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury. Zmierz stopę odzyskiwania informacji na ustalonej serii pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko naprawiają słabe mechanizmy wyszukiwania.
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,
)
Test rozmowy wieloetapowej
Gdy pracujesz nad testem rozmowy wieloetapowej, najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Zdokumentuj zarówno ścieżkę prawidłowego działania, jak i ścieżkę naprawczą. Próby ponowne, kontrola przez ludzi oraz obsługa wiadomości nieodebranych stanowią część produktu, a nie elementy dopinane później. Zmierz stopień odzyskiwania informacji na ustalonej serii pytań przed dostosowywaniem promptów. Częste zmiany promptów rzadko naprawiają słabe mechanizmy wyszukiwania informacji.
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)
Gdy pracujesz nad testem rozmowy wieloetapowej, najpierw zapisz umowę: wymagane dane wejściowe, sygnał sukcesu oraz to, co dzieje się w przypadku częściowego niepowodzenia. Taka lista kontrolna zapewnia uczciwość późniejszych zmian w kodzie. Traktuj ten etap jako umowę pomiędzy danymi wejściowymi a zweryfikowanymi wynikami. Nazwij poszczególne elementy, zdefiniuj kryteria sukcesu i odrzuć ciche, częściowe ukończenie zadań.
Krok 6: Kod agenta produkcyjnego (agent.py)
Krok 6: Kod agenta produkcyjnego (agent.py) funkcjonuje najlepiej, gdy traktowany jest jako mierzalna powierzchnia do analizy. Zapisz jeden idealny przykład działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres pracy.
Zapisuj czasy wykonywania oraz koszt tokenów lub zapytań obok wyników funkcjonalnych. Wczesna widoczność kosztów zapobiega nieoczekiwanym rachunkom, gdy przechodzi się z środowiska demonstracyjnego do współdzielonych środowisk.
Rozdziel politykę dzielenia na fragmenty od polityki wyszukiwania. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
# 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)
Konfiguracja (agent-config.yaml)
Konfiguracja (agent-config.yaml) funkcjonuje najlepiej, gdy jest traktowana jako coś mierzalnego. Zapisz jeden idealny przykład działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres.
Trzymaj konfigurację poza kodem aplikacji. Pliki środowiskowe, magazyny tajnych danych oraz flagi funkcjonalne powinny znajdować się w jednym miejscu, które operatorzy mogą sprawdzić bez konieczności czytania całej struktury.
Oddziel zasadę dzielenia na fragmenty od zasady pobierania danych. Zmiana jednej z nich nie powinna zmuszać do przepisywania drugiej, gdy zmieniają się metryki jakości.
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>
Krok 7: Rejestracja, logowanie i wdrożenie
Krok 7: Rejestracja, logowanie i wdrożenie funkcjonuje najlepiej, gdy jest traktowany jako coś mierzalnego. Zapisz jeden idealny przykład działania, jeden przypadek awarii oraz notatkę dotyczącą cofnięcia zmian, zanim rozszerzysz zakres.
Logowanie do 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
Rejestracja w 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}")
Wdrożenie
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}")
Rezultat
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}}
)
Szybka weryfikacja trwałości pamięci:
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."
)