Ваш агент 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 работает наилучшим образом, когда его рассматривают как измеримую систему. Соберите один идеальный пример работы, один случай сбоя и записку о возврате к предыдущему состоянию перед расширением объема работ. Документируйте как успешный, так и восстановительный пути выполнения. Повторные попытки, проверки человеком и обработка неработающих сообщений являются частью продукта, а не этапом последующей доработки. Разделяйте политику разбиения данных на части и политику их извлечения. Изменение одной из них не должно приводить к переписыванию другой при изменении показателей качества. Проект автомасштабирования Lakebase работает наилучшим образом, когда его рассматривают как измеримую систему. Соберите один идеальный пример работы, один случай сбоя и записку о возврате к предыдущему состоянию перед расширением объема работ. Рассматривайте этот этап как контракт между входными данными и проверенными выходными результатами. Дайте названия элементам документации, определите критерии успеха и не соглашайтесь на молчаливое частичное выполнение задач.
Получение деталей подключения программным способом
Чтобы получить детали подключения программно, необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Зафиксируйте время выполнения, а также стоимость токенов или запросов рядом с функциональными результатами. Отображение стоимости заранее помогает избежать неожиданных счетов при переходе с демо-среды в общедоступные среды. Укажите те части текста, которые фактически легли в основу ответа. Без цитат операторы не смогут отличить галлюцинации от пробелов в индексации.
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: Создание агента, учитывающего контекст, сначала запишите спецификации: необходимые входные данные, сигнал успешного выполнения и действия при частичной неудаче. Такой чек-лист поможет сохранять честность при последующих изменениях кода. Записывайте время выполнения и стоимость токенов или запросов рядом с функциональными результатами. Отслеживание затрат на раннем этапе предотвращает неожиданные счета при переходе от демо-среды к общедоступным средам. Оцените уровень воспроизводимости ответов на фиксированный набор вопросов перед настройкой подсказок. Частая смена подсказок редко помогает улучшить качество поиска информации.
Интерактивный агент (ноутбук)
При работе с интерактивным агентом (ноутбуком) сначала запишите условия работы: необходимые входные данные, сигнал успешного выполнения и действия при частичной неудаче. Такой список помогает сохранять честность при последующих изменениях кода. Храните конфигурацию вне кода приложения. Файлы среды, хранилища секретов и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая весь код. Измеряйте уровень воспроизводимости ответов на фиксированном наборе вопросов перед настройкой подсказок. Частая смена подсказок редко помогает улучшить качество поиска информации.
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."
)