Ваш RAG-агент забуває все після одного повідомлення – ось як ви це виправили за допомогою Databricks…
Покрокова інструкція щодо вирішення проблеми: ваш RAG-агент забуває все після одного повідомлення – ось як ви це виправили за допомогою Databricks…: контракти, перевірки та готові блоки коду для команд.
У цьому посібнику описано процес створення системи від сировини до готового продукту для ситуації, коли ваш RAG-агент забуває все після одного повідомлення – ось як я це виправив за допомогою Databricks Lakebase. Основна увага приділяється конкретним крокам виконання, чітким перевіркам та коду, який можна просто додати до репозиторію без необхідності здогадуватися щодо його призначення. Для загального огляду спочатку визначте вхідні дані, виконавця кроку та критерії завершення перед зміною коду. Виконавці мають мати можливість перезапустити крок з відомої точки контролю, не намагаючись здогадатися про прихований стан системи. Розглядайте цю стадію як контракт між вхідними даними та перевіреними результатами. Призначте назви файлів, визначте критерії успіху та не допускайте безповідомного часткового завершення роботи.
Архітектура
Під час роботи над «Архітектурою» спочатку запишіть умови використання: необхідні вхідні дані, сигнал про успішне виконання та те, що відбувається у разі часткової невдачі. Такий перелік допомагає зберігати чесність пізніших змін у коді. Запишіть час виконання та витрати на токени або запити поруч із результатами функціоналу. Чітке бачення витрат заздалегідь запобігає несподіваним рахункам, коли система переходить від демо-режиму до спільних середовищ. Вимірюйте рівень точності відповідей на фіксований набір запитань перед налаштуванням промптів. Часта зміна промптів рідко допомагає покращити якість пошуку.
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)
Крок 1: Аналіз документів за допомогою ai_parse_document()
Під час виконання Кроку 1: Аналіз документів за допомогою ai_parse_document(), спочатку запишіть умови використання: необхідні вхідні дані, сигнал про успіх та наслідки часткової невдачі. Такий перелік допоможе зберегти чесність подальших змін у коді.
Зберігайте конфігурацію окремо від коду додатку. Файли середовища, сховища конфіденційних даних та флаги функціоналу мають знаходитися в одному місці, щоб оператори могли їх перевіряти, не читаючи весь код.
Перед налаштуванням запитів вимірюйте рівень відтворення інформації на фіксованому наборі запитань. Часта зміна формулювань запитів рідко допомагає покращити ефективність пошуку.
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}")
Крок 2: Очищення, трансформація та розділення на частини
Під час виконання Кроку 2: Очищення, трансформація та розділення на частини, спочатку запишіть умови взаємодії: необхідні вхідні дані, сигнал про успіх та те, що відбувається у разі часткової невдачі. Цей перелік допомагає зберігати чесність подальших змін у коді. Документуйте як шлях успішної роботи, так і шлях відновлення. Повторні спроби, людський контроль та обробка некоректних повідомлень є частиною продукту, а не етапом подальшої оптимізації. Вимірюйте рівень відтворення інформації за фіксованим набором запитань перед налаштуванням запрошень до введення даних. Зміна запрошень рідко вирішує проблеми слабкого пошуку інформації. Під час виконання Кроку 2: Очищення, трансформація та розділення на частини, спочатку запишіть умови взаємодії: необхідні вхідні дані, сигнал про успіх та те, що відбувається у разі часткової невдачі. Цей перелік допомагає зберігати чесність подальших змін у коді. Розглядайте цей етап як умови взаємодії між вхідними даними та перевіреними результатами. Назвіть створювані елементи, визначте критерії успіху та не допускайте беззвучного часткового виконання завдань.
Швидке вилучення простого тексту
Швидке вилучення простого тексту найкраще функціонує, якщо його розглядати як вимірювану поверхню. Збережіть один ідеальний запис, один випадок збою та примітку про скасування змін перед розширенням обсягу роботи. Записуйте час виконання та витрати на обробку токенів чи запитів поруч із функціональними результатами. Візуалізація витрат на ранньому етапі запобігає несподіваним рахункам під час переходу від демо-середовища до спільних середовищ. Розділяйте політику розбиття на частини та політику пошуку. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
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
)
Розбиття на частини за допомогою LangChain
Чанк з LangChain працює найкраще, коли його розглядають як вимірювану поверхню. Збережіть один ідеальний запис, один випадок збою та примітку про скасування змін перед розширенням обсягу роботи. Тримайте конфігурацію поза кодом додатку. Файли середовища, сховища конфіденційних даних та флаги функцій мають знаходитися в одному місці, де оператори можуть їх перевіряти, не читаючи весь граф. Розділяйте політику формування чанків та політику пошуку. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
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)
Крок 3: Створення векторного пошуку
Крок 3: Механізм пошуку векторів працює найкраще, коли його розглядають як вимірювану поверхню. Зафіксуйте один ідеальний варіант роботи, один випадок збою та примітки щодо скасування змін перед розширенням обсягу. Документуйте як успішний, так і відновлювальний сценарії роботи разом. Повторні спроби, людський контроль та обробка некоректних повідомлень є частиною продукту, а не етапом подальшої оптимізації. Розділіть політику часткової обробки даних від політики їх пошуку. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
Увімкніть подачу даних про зміни
Функція Change Data Feed працює найкраще, коли її розглядають як вимірювану систему. Збережіть один ідеальний зразок даних, один випадок збою та примітки щодо скасування змін перед розширенням обсягу роботи. Віддавайте перевагу невеликим, тестованим одиницям перед складними скриптами. Коли якийсь крок зазнає невдачі, причина має вказувати на конкретну відповідальність, а не на заплутану послідовність дій. Розділяйте політику часткової обробки даних та політику їх отримання. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
ALTER TABLE <YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
Створення індексу синхронізації Delta
Створення індексу Delta Sync працює найкраще, якщо його розглядати як вимірювану поверхню. Збережіть один ідеальний протокол, один випадок збою та запис про скасування змін перед розширенням обсягу роботи. Розглядайте цей етап як контракт між вхідними даними та перевіреними результатами. Позначте всі елементи, визначте критерії успіху та не допускайте мовчазного часткового виконання завдань. Розділіть політику часткового оброблення даних від політики їх отримання. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
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",
)
Тестування отримання даних
Тестування отримання даних працює найкраще, коли його розглядають як вимірювану характеристику. Запишіть один ідеальний приклад обробки даних, один випадок невдачі та примітки щодо скасування змін перед розширенням обсягу роботи. Записуйте час виконання та витрати на токени чи запити поруч із функціональними результатами. Візуалізація витрат заздалегідь запобігає несподіваним рахункам під час переходу від демо-середовища до спільних середовищ. Розділяйте політику часткової обробки даних та політику їх отримання. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
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)
Крок 4: Налаштування Lakebase для пам’яті розмов
Крок 4: Налаштування Lakebase для пам’яті розмов працює найкраще, коли її розглядають як вимірювану поверхню. Збережіть один ідеальний запис розмови, один випадок збою та примітку про скасування змін перед розширенням обсягу функціоналу. Тримайте конфігурацію окремо від коду додатку. Файли середовища, сховища конфіденційних даних та флаги функцій мають знаходитися в одному місці, де оператори можуть їх перевіряти, не читаючи весь код. Розділіть політику часткового оброблення даних від політики їх пошуку. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
Створення проекту автомасштабування Lakebase
Проєкт Lakebase Autoscaling працює найкраще, коли його розглядають як вимірювану систему. Збережіть один ідеальний запис, один випадок збою та примітку щодо скасування змін перед розширенням обсягу роботи. Документуйте як успішний, так і відновлювальний сценарії роботи. Повторні спроби, людський контроль та обробка некоректних повідомлень є частиною продукту, а не етапом подальшої оптимізації. Розділіть політику часткової обробки даних від політики їх отримання. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості. Проєкт Lakebase Autoscaling працює найкраще, коли його розглядають як вимірювану систему. Збережіть один ідеальний запис, один випадок збою та примітку щодо скасування змін перед розширенням обсягу роботи. Розглядайте цей етап як контракт між вхідними даними та перевіреними результатами. Позначте всі елементи, визначте критерії успіху та не допускайте мовчазного часткового виконання завдань.
Отримати деталі підключення програмно
Щоб отримати деталі з’єднання програмно, необхідно визначити вхідні дані, власника кроку та критерії завершення перед зміною коду. Оператори повинні мати можливість перезапустити крок з відомої точки контролю, не намагаючись вгадати прихований стан. Записуйте час виконання та витрати на токени або запити поруч із функціональними результатами. Візуалізація витрат заздалегідь запобігає несподіваним рахункам під час переходу з демо-середовища у спільні. Наводьте ті частини тексту, які фактично лягли в основу відповіді. Без посилань оператори не зможуть відрізнити галюцинації від проблем з індексуванням.
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}")
Перевірка з’єднання
Щоб перевірити з’єднання, необхідно визначити вхідні дані, власника кроку та критерії завершення перед зміною коду. Оператори повинні мати можливість знову запустити крок з відомої точки контролю, не намагаючись вгадати прихований стан. Конфігурацію слід тримати окремо від коду додатку. Файли середовища, сховища конфіденційних даних та флаги функцій мають знаходитися в одному місці, де оператори можуть їх перевірити, не читаючи весь код. Наводьте уривки тексту, які фактично лежать в основі відповіді. Без посилань оператори не зможуть відрізнити галюцинацію від проблем з індексуванням.
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!")
Створення таблиць точок контролю
Для створення таблиць контрольних точок необхідно визначити вхідні дані, власника кроку та критерії завершення ще до зміни коду. Оператори повинні мати можливість знову виконати крок з відомої контрольної точки, не намагаючись вгадати прихований стан. Необхідно документувати як шлях успішного виконання, так і шлях відновлення. Повторні спроби, людський контроль та обробка некоректних повідомлень є частиною продукту, а не етапом подальшої оптимізації. Наводьте ті уривки, які фактично лежать в основі відповіді. Без посилань оператори не зможуть відрізнити галюцинації від проблем з індексуванням.
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!")
Для створення таблиць контрольних точок необхідно визначити вхідні дані, власника кроку та критерії завершення ще до зміни коду. Оператори повинні мати можливість знову виконати крок з відомої контрольної точки, не намагаючись вгадати прихований стан. Розглядайте цей етап як контракт між вхідними даними та перевіреними результатами. Позначте всі елементи, визначте критерії успіху та не допускайте мовчазного часткового виконання завдань.
Крок 5: Створення агента, що розуміє контекст
Під час виконання Кроку 5: Створення агента, що розуміє контекст, спочатку запишіть умови роботи: необхідні вхідні дані, сигнал про успіх та те, що відбувається у разі часткової невдачі. Цей перелік допомагає зберігати чесність пізніших змін у коді. Записуйте час виконання та витрати на токени або запити поруч із функціональними результатами. Відображення витрат заздалегідь запобігає несподіваним рахункам, коли система переходить від демо-режиму до спільних середовищ. Перевіряйте рівень відтворення інформації на фіксованому наборі запитань перед налаштуванням підказок. Часта зміна підказок рідко допомагає покращити ефективність пошуку інформації.
Інтерактивний агент (ноутбук)
Під час роботи з Interactive Agent (Notebook) спочатку запишіть контракт: необхідні вхідні дані, сигнал про успіх та те, що відбувається при частковій невдачі. Такий перелік допомагає зберігати чесність пізніших змін у коді. Зберігайте конфігурацію окремо від коду додатку. Файли середовища, сховища секретних даних та флаги функцій мають знаходитися в одному місці, де оператори можуть їх перевіряти, не читаючи весь код. Вимірюйте рівень відтворення інформації на фіксованому наборі запитань перед налаштуванням підказок. Часта зміна підказок рідко виправляє проблеми з недостатньою ефективністю пошуку.
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,
)
Тестування багатокрокової розмови
Під час роботи над тестом багатокрокової розмови спочатку запишіть угоду: необхідні вхідні дані, сигнал про успіх та те, що відбувається при частковій невдачі. Цей перелік допомагає зберігати чесність у подальших змінах коду. Одночасно задокументуйте шлях успішного виконання та шлях відновлення. Повторні спроби, людський контроль та обробка некоректних повідомлень є частиною продукту, а не етапом подальшої доробки. Вимірюйте рівень відтворення інформації за фіксованим набором запитань перед налаштуванням підказок. Зміна підказок рідко вирішує проблеми слабкого пошуку інформації.
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)
Під час роботи над тестом багатокрокової розмови спочатку запишіть угоду: необхідні вхідні дані, сигнал про успіх та те, що відбувається при частковій невдачі. Цей перелік допомагає зберігати чесність у подальших змінах коду. Розглядайте цей етап як угоду між вхідними даними та перевіреними результатами. Позначте всі елементи, визначте критерії успіху та не допускайте мовчазного часткового виконання завдань.
Крок 6: Код агента для продакшну (agent.py)
Крок 6: Код агента для продакшну (agent.py) найкраще функціонує, якщо його розглядати як вимірювану систему. Збережіть один ідеальний запис роботи, один випадок збою та примітки щодо скасування змін перед розширенням обсягу завдань.
Записуйте час виконання та витрати на токени чи запити поруч із функціональними результатами. Відображення витрат на ранньому етапі запобігає несподіваним рахункам під час переходу від демо-середовища до спільних середовищ.
Розділяйте політику часткового оброблення даних та політику їх пошуку. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
# 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)
Конфігурація (agent-config.yaml)
Конфігурація (agent-config.yaml) функціонує найкраще, якщо її розглядати як показник для вимірювання. Збережіть один ідеальний зразок запису, один випадок збою та примітку щодо скасування змін перед розширенням обсягу роботи.
Тримайте конфігурацію окремо від коду додатку. Файли середовища, сховища конфіденційних даних та флаги функцій мають знаходитися в одному місці, де оператори можуть їх перевіряти, не читаючи весь код.
Розділіть політику часткового оброблення даних від політики їх отримання. Зміна однієї з них не повинна змушувати переписувати іншу при зміні показників якості.
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>
Крок 7: Логування, реєстрація та розгортання
Крок 7: Логування, реєстрація та розгортання функціонує найкраще, якщо його розглядати як показник для вимірювання. Збережіть один ідеальний зразок запису, один випадок збою та примітку щодо скасування змін перед розширенням обсягу роботи.
Логування в 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
Реєстрація в 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}")
Розгортання
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}")
Результат
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}}
)
Швидка перевірка на наявність стабільності пам’яті:
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."
)