Ваш агент 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: Чыстка, трансфармація і разбіўка на часткі, спачатку запісайце угоду: неабяжлівыя данні, сигнал успеху і тое, што выходзіць у разе частковага невыпання. Такі список пераконтроўкаў дапамагае заліцьваты пазнейшыя змены ў кодзе.
Шыраная выкарыстоўвання простага тексту
Шырокая выкарыстоўка простага тэксту работае наяўней, калі яе спрыявае меркаванне на аснове певных параметраў. Перш чым расширваць масштабы, зафіксавайце адны ідеальны прыклад, адзін прыклад неудачы і запіс пра можлівасць вярнуцься да пачатковага стану. Запісвайце часы выконання, а таксу на обробку токеноў чы запытак праза функцыйнае рэзультат. Відчутнасць костаў з самага пачатку запобегае неспакойным рахункам, калі процес пераходзіць з дэмовай среды ў спакульную. Раздзеляйце правілы часткавання інфармацыі ад правіл яе выкарыстоўкі. Змена ў адных не павинна вымагаць перапісвы іншых, калі зменяюцца паметры якосці.
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
Функцыя Change Data Feed працюе наўсёй краща, калі яе розглядаць як мерыемую структуру. Зафіксавайце адну ідеальную версію дадзеных, адны прыклад неудачы і прыметкі па поверненню да пярвоначальнага стану, прычаму расшырэння сферы дзеяння. Валідзіце маленькія, тэставаныя елементы замест амаль неконтрольваных скрыптав. Калі якісь крок не выйшае, прычына неудачы павінна вказываць на адную конкрэтную адпаведальнасць, а не на заплутаны ланцюг задач. Раздзеліце правілы часткавання дадзеных ад правіл ўтрымання іх. Змена адных не павінна вымагаць перапісву іншых, калі зменяюцыся паказателі якосці.
ALTER TABLE <YOUR_CATALOG>.<YOUR_SCHEMA>.docs_chunked
SET TBLPROPERTIES (delta.enableChangeDataFeed = true);
Створыце індэкс Delta Sync
Процес стварэння індекса 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."
)