4
2

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Lakebase × LangGraph × Qwen3で作る、セッションを跨ぐ記憶を持つAIエージェント (連載 第2回)

4
Last updated at Posted at 2026-04-28

はじめに

前回の記事で、Databricks Lakebase Autoscaling と LangGraph の PostgresSaver (Checkpointer) を使って、短期メモリ付きのエージェントを構築しました。会話の文脈を thread_id 単位で保持できるようになり、同一スレッド内では「私の名前は太郎です」と伝えた直後に「私の名前は何でしたっけ?」と聞いても、ちゃんと「太郎さん」と返してくれるようになりました。

ただ、これには明確な限界があります。thread_id を変える(別セッションを始める)と、エージェントは前の会話を一切覚えていません。これでは「数日前に話した好み」「先週の打ち合わせの内容」「ユーザーが過去に教えてくれた名前」などをエージェントに引き継がせることはできません。

これを解決するのが 長期メモリ です。LangGraph の PostgresStore (Store API) を使うと、thread_id を跨いで永続化される、ユーザー単位のメモリストアを実装できます。本記事では前回のエージェントを拡張して、セッション横断で記憶を引き継ぐエージェントを構築します。

本記事は連載の第2回です。

短期メモリと長期メモリの違い

LangGraph には2つの永続化レイヤーがあり、それぞれ別の役割を担います。

観点 Checkpointer (短期メモリ) Store (長期メモリ)
実装クラス PostgresSaver PostgresStore
スコープ thread_id 単位 任意 (user_id, org_id など namespace で自由)
保存内容 会話の状態スナップショット (messages, ツール呼び出し履歴など) 任意の構造化データ (JSON) + 意味検索用 Embedding
保存タイミング エージェントの各ステップで自動 ツール経由 or アプリのコードで明示的に
検索方法 thread_id をキーに最新状態を取得 namespace + キー、または意味検索 (ベクトル類似度)
想定される用途 「前のターンで何を言ったか」 「過去にユーザーが教えてくれた好み・属性」

両者は 独立して動作 します。Checkpointer だけ、Store だけ、両方、のどの組み合わせでも構築可能です。本記事では両方を併用したエージェントを作ります。

構築するもの

前回作ったエージェントを拡張する形で進めます。具体的には以下のツールを追加します。

  • save_memory(content): ユーザーに関する情報を Store に保存
  • recall_memories(query): 過去のメモリを意味検索で参照

エージェントは system_prompt で「ユーザーが自分について話してくれたら save_memory、思い出す必要があれば recall_memories を呼ぶ」と誘導します。

完成形のアーキテクチャ:

[ユーザー]
   ↓
[LangGraph Agent (create_react_agent)]
   ├─ ChatDatabricks (Foundation Model API)
   ├─ Tools: get_current_datetime, save_memory, recall_memories
   ├─ PostgresSaver (Checkpointer) ──┐
   └─ PostgresStore (Store) ─────────┤
                                     ↓
                          [Lakebase Autoscaling]
                       checkpoints テーブル群: thread_id 単位
                       store / store_vectors: namespace 単位
                       (pgvector + Qwen3-Embedding でベクトル検索)

前提と動作確認バージョン

  • 第1回の構成 (Lakebase Autoscaling Project + 接続成功) が動いていること
  • databricks-qwen3-embedding-0-6b Foundation Model API エンドポイントが利用可能なこと

動作確認時のバージョン:

psycopg: 3.3.3
psycopg-binary: 3.3.3
psycopg-pool: 3.3.0
langgraph: 1.1.10
langgraph-checkpoint: 4.0.3
langgraph-checkpoint-postgres: 3.0.5
databricks-sdk: 0.105.0
databricks-langchain: 0.19.0
PostgreSQL: 17.8 (Lakebase 側)
pgvector: 0.8.0
Embedding model: databricks-qwen3-embedding-0-6b (1024次元)
LLM: databricks-claude-sonnet-4

ステップ1: 環境準備

新規ノートブックを作成し、サーバレスコンピュートにアタッチします。第1回と同じパッケージを再インストールします。

%pip install "psycopg[binary,pool]" databricks-sdk
%restart_python
import psycopg
import databricks.sdk

print(f"psycopg: {psycopg.__version__}")
print(f"psycopg impl: {psycopg.pq.__impl__}")
print(f"libpq: {psycopg.pq.version()}")
print(f"databricks-sdk: {databricks.sdk.version.__version__}")

databricks-sdk が 0.94.0 未満の場合、w.postgres API が使えないので明示的にアップグレードします。

%pip install --upgrade "databricks-sdk>=0.94.0"
%restart_python

ステップ2: Lakebase接続 + Checkpointerのセットアップ

第1回と同じ手順で接続します。第1回で作成した agent-memory-db プロジェクトをそのまま使います。

%pip install langgraph langgraph-checkpoint-postgres
%restart_python
import concurrent.futures
from databricks.sdk import WorkspaceClient
from psycopg_pool import ConnectionPool
from langgraph.checkpoint.postgres import PostgresSaver

w = WorkspaceClient()

PROJECT_ID = "agent-memory-db"
BRANCH_ID = "production"
ENDPOINT_ID = "primary"
HOST = "<your endpoint host>"

cred = w.postgres.generate_database_credential(
    endpoint=f"projects/{PROJECT_ID}/branches/{BRANCH_ID}/endpoints/{ENDPOINT_ID}"
)
USER = w.current_user.me().user_name

conninfo = (
    f"dbname=databricks_postgres user={USER} password={cred.token} "
    f"host={HOST} port=5432 sslmode=require"
)

def init_pool_and_saver():
    pool = ConnectionPool(
        conninfo=conninfo,
        max_size=5,
        kwargs={"autocommit": True, "prepare_threshold": 0},
        open=True,
    )
    checkpointer = PostgresSaver(pool)
    checkpointer.setup()  # 第1回でテーブルが作成済みなら冪等にスキップされる
    return pool, checkpointer

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    pool, checkpointer = ex.submit(init_pool_and_saver).result(timeout=60)

print("Pool and checkpointer ready")

第1回で作成したテーブルとチェックポイントがそのまま残っていることを確認できます。

def list_tables_and_count():
    with pool.connection() as conn:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT tablename FROM pg_tables 
                WHERE schemaname = 'public' 
                ORDER BY tablename;
            """)
            tables = cur.fetchall()
            cur.execute("SELECT COUNT(*) FROM checkpoints;")
            ckpt_count = cur.fetchone()[0]
            return tables, ckpt_count

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    tables, ckpt_count = ex.submit(list_tables_and_count).result(timeout=15)

print("Tables in public schema:")
for t in tables:
    print(f"  - {t[0]}")
print(f"\nExisting checkpoints: {ckpt_count}")

出力:

Tables in public schema:
  - checkpoint_blobs
  - checkpoint_migrations
  - checkpoint_writes
  - checkpoints

Existing checkpoints: 9

第1回で tutorial-session-001tutorial-session-002 で作った 9 件のチェックポイントがそのまま残っていることが確認できました。Lakebase Autoscaling のデータ永続性が体感できる場面です。

ステップ3: pgvector有効化とEmbeddingエンドポイントの動作確認

長期メモリで意味検索を使うために、pgvector 拡張を有効化します。Lakebase は PostgreSQL ベースで pgvector がサポートされています。

import concurrent.futures

def check_extensions():
    with pool.connection() as conn:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT name, default_version, installed_version
                FROM pg_available_extensions
                WHERE name = 'vector';
            """)
            return cur.fetchone()

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    result = ex.submit(check_extensions).result(timeout=15)

print(f"vector extension: {result}")

出力例:

vector extension: ('vector', '0.8.0', None)

installed_versionNone なら、まだインストールされていない (CREATE EXTENSION 未実行) という意味です。有効化します。

def enable_vector_extension():
    with pool.connection() as conn:
        with conn.cursor() as cur:
            cur.execute("CREATE EXTENSION IF NOT EXISTS vector;")
        with conn.cursor() as cur:
            cur.execute("SELECT extname, extversion FROM pg_extension WHERE extname = 'vector';")
            return cur.fetchone()

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    result = ex.submit(enable_vector_extension).result(timeout=15)

print(f"vector extension: {result}")

出力:

vector extension: ('vector', '0.8.0')

次に Embedding モデルを選びます。本記事では Qwen3-Embedding-0.6B を使います。理由は以下の通りです。

  • 多言語対応が強く、日本語の意味理解が特に良い
  • Foundation Model API のマネージドエンドポイントとして提供されている (databricks-qwen3-embedding-0-6b)
  • 出力次元は 1024 (Matryoshka 表現でカスタマイズも可能)

動作確認します。

def test_embedding(endpoint_name: str):
    response = w.serving_endpoints.query(
        name=endpoint_name,
        input=["こんにちは、世界。"],
    )
    embedding = response.data[0].embedding
    return len(embedding), embedding[:3]

dim, sample = test_embedding("databricks-qwen3-embedding-0-6b")
print(f"databricks-qwen3-embedding-0-6b: dim={dim}, sample={sample}")

出力例:

databricks-qwen3-embedding-0-6b: dim=1024, sample=[-0.00996891874819994, -0.015331508591771126, -0.005946975201368332]

1024次元のベクトルが返ってきており、日本語入力に対しても問題なく動作しています。

ステップ4: PostgresStore のセットアップ

Embedding 関数を作って、PostgresStore に渡します。ここで一つ重要な注意点があります。Qwen3-Embedding が稀に整数値 (例えば 0) を含むベクトルを返すことがあり、psycopg + pgvector に書き込むときに「mixed types」エラーになります。これを防ぐため、明示的に float() でキャストします。

def embed_texts(texts: list[str]) -> list[list[float]]:
    """Databricks の Qwen3-Embedding を呼び出して埋め込みベクトルを返す。
    
    Note: Qwen3 が稀に int 値を含むベクトルを返すため、
    pgvector に渡す前に明示的に float にキャストする。
    """
    response = w.serving_endpoints.query(
        name="databricks-qwen3-embedding-0-6b",
        input=texts,
    )
    return [[float(x) for x in item.embedding] for item in response.data]

# 動作確認
test_vec = embed_texts(["温泉が好きです"])[0]
print(f"dim: {len(test_vec)}")
print(f"all float: {all(isinstance(x, float) for x in test_vec)}")
print(f"sample: {test_vec[:5]}")

all float: True を確認できたら、PostgresStore をセットアップします。

import concurrent.futures
from langgraph.store.postgres import PostgresStore

def setup_store():
    store = PostgresStore(
        pool,
        index={
            "dims": 1024,
            "embed": embed_texts,
        },
    )
    store.setup()
    return store

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    store = ex.submit(setup_store).result(timeout=60)

print("PostgresStore setup OK")

store.setup() は内部で以下のテーブルを Lakebase 側に作成します。

  • store: メモリ本体 (prefix, key, value JSONB, タイムスタンプ, TTL)
  • store_vectors: pgvector を使った Embedding 保存テーブル
  • store_migrations, vector_migrations: スキーマバージョン管理

テーブル構成を確認します。

def list_all_tables():
    with pool.connection() as conn:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT tablename FROM pg_tables 
                WHERE schemaname = 'public' 
                ORDER BY tablename;
            """)
            return cur.fetchall()

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    tables = ex.submit(list_all_tables).result(timeout=15)

print("Tables in public schema:")
for t in tables:
    print(f"  - {t[0]}")

出力:

Tables in public schema:
  - checkpoint_blobs
  - checkpoint_migrations
  - checkpoint_writes
  - checkpoints
  - store
  - store_migrations
  - store_vectors
  - vector_migrations

第1回からの Checkpointer 系 4テーブル + 今回追加した Store 系 4テーブルで、計8テーブルになりました。

Storeの基本動作確認 (エージェントを介さず素のAPIで)

エージェントに組み込む前に、Store単体の動作を確認しておきます。

import concurrent.futures

def test_store_basic():
    namespace = ("memories", "user-taro")
    
    # 保存
    store.put(
        namespace=namespace,
        key="memory-001",
        value={"text": "私は温泉が好きで、特に箱根によく行きます"},
    )
    
    # キーで取得
    item = store.get(namespace=namespace, key="memory-001")
    
    # 意味検索
    results = store.search(
        namespace,
        query="リラックスできる場所",
        limit=3,
    )
    
    return item, results

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    item, results = ex.submit(test_store_basic).result(timeout=30)

print("=== Direct get ===")
print(item)
print("\n=== Semantic search for 'リラックスできる場所' ===")
for r in results:
    print(f"  score={r.score:.4f}  value={r.value}")

出力:

=== Direct get ===
Item(namespace=['memories', 'user-taro'], key='memory-001', value={'text': '私は温泉が好きで、特に箱根によく行きます'}, ...)

=== Semantic search for 'リラックスできる場所' ===
  score=0.4577  value={'text': '私は温泉が好きで、特に箱根によく行きます'}

「温泉が好き」というメモリに対して、「リラックスできる場所」という直接の単語マッチではないクエリが意味検索で引っかかっています。スコア 0.4577 は cosine similarity ベースで、Qwen3 が「温泉」と「リラックスできる場所」を意味的に近いと判定できている証拠です。

ステップ5: メモリ操作ツールとエージェント拡張

エージェント側のパッケージを追加します。

%pip install databricks-langchain
%restart_python

第1回と同じく、databricks-langchain をインストールすると langgraph がダウングレードされて ImportError: cannot import name 'ExecutionInfo' が出ることがあります。バージョン確認の上、必要に応じてアップグレードします。

import subprocess
result = subprocess.run(["pip", "list", "--format=freeze"], capture_output=True, text=True)
for line in result.stdout.split("\n"):
    if any(k in line.lower() for k in ["psycopg", "langgraph", "langchain", "databricks-"]):
        print(line)

langgraph==1.0.x になっていたら:

%pip install --upgrade "langgraph>=1.1.10" "langgraph-checkpoint-postgres"
%restart_python

%restart_python で状態がリセットされたので、Pool / Checkpointer / Store を再構築します。

import concurrent.futures
from databricks.sdk import WorkspaceClient
from psycopg_pool import ConnectionPool
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.store.postgres import PostgresStore

w = WorkspaceClient()

PROJECT_ID = "agent-memory-db"
BRANCH_ID = "production"
ENDPOINT_ID = "primary"
HOST = "<your endpoint host>"

cred = w.postgres.generate_database_credential(
    endpoint=f"projects/{PROJECT_ID}/branches/{BRANCH_ID}/endpoints/{ENDPOINT_ID}"
)
USER = w.current_user.me().user_name

conninfo = (
    f"dbname=databricks_postgres user={USER} password={cred.token} "
    f"host={HOST} port=5432 sslmode=require"
)

def embed_texts(texts: list[str]) -> list[list[float]]:
    response = w.serving_endpoints.query(
        name="databricks-qwen3-embedding-0-6b",
        input=texts,
    )
    return [[float(x) for x in item.embedding] for item in response.data]

def init_all():
    pool = ConnectionPool(
        conninfo=conninfo,
        max_size=5,
        kwargs={"autocommit": True, "prepare_threshold": 0},
        open=True,
    )
    checkpointer = PostgresSaver(pool)
    store = PostgresStore(
        pool,
        index={"dims": 1024, "embed": embed_texts},
    )
    return pool, checkpointer, store

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    pool, checkpointer, store = ex.submit(init_all).result(timeout=60)

print("Pool, checkpointer, and store ready")

メモリ保存・参照ツールを定義し、エージェントに組み込みます。

from datetime import datetime
from typing import Annotated
from langchain_core.tools import tool
from langgraph.prebuilt import create_react_agent, InjectedStore
from langgraph.store.base import BaseStore
from langgraph.config import get_config
from databricks_langchain import ChatDatabricks


@tool
def get_current_datetime() -> str:
    """現在の日時を返します。"""
    return datetime.now().strftime("%Y-%m-%d %H:%M:%S")


@tool
def save_memory(
    content: str,
    store: Annotated[BaseStore, InjectedStore()],
) -> str:
    """ユーザーに関する長期記憶として情報を保存します。
    
    ユーザーの好み、特徴、過去の発言などで、将来のセッションでも覚えておくべき情報を保存してください。
    例: 趣味、好きな食べ物、家族構成、職業など。
    """
    config = get_config()
    user_id = config["configurable"].get("user_id", "default")
    namespace = ("memories", user_id)
    
    key = f"memory-{datetime.now().strftime('%Y%m%d-%H%M%S-%f')}"
    
    store.put(
        namespace=namespace,
        key=key,
        value={"text": content, "saved_at": datetime.now().isoformat()},
    )
    return f"記憶しました: {content}"


@tool
def recall_memories(
    query: str,
    store: Annotated[BaseStore, InjectedStore()],
) -> str:
    """ユーザーに関する過去の記憶を意味検索で参照します。
    
    ユーザーの好みや過去の発言を思い出したい時に使ってください。
    クエリには思い出したい内容を自然言語で書いてください。
    """
    config = get_config()
    user_id = config["configurable"].get("user_id", "default")
    namespace = ("memories", user_id)
    
    results = store.search(namespace, query=query, limit=5)
    
    if not results:
        return "該当する記憶はありませんでした。"
    
    lines = ["過去の記憶 (関連度順):"]
    for r in results:
        lines.append(f"  - {r.value.get('text', '')} (関連度: {r.score:.3f})")
    return "\n".join(lines)


llm = ChatDatabricks(endpoint="databricks-claude-sonnet-4")

system_prompt = """あなたはユーザーとの長期的な関係を築くアシスタントです。

ユーザーが自分自身について何か新しいことを話してくれたら、save_memory ツールで保存してください。
ユーザーから何か質問されたとき、関連する過去の情報がありそうなら recall_memories ツールで思い出してください。
ユーザーの言葉や好みを尊重し、覚えていることを自然に会話に織り込んでください。"""

agent = create_react_agent(
    model=llm,
    tools=[get_current_datetime, save_memory, recall_memories],
    checkpointer=checkpointer,
    store=store,
    prompt=system_prompt,
)

print("Agent with long-term memory ready")

ポイントは以下の3つです。

  • create_react_agentstore=store を渡すことで、ツールから InjectedStore 経由で Store にアクセスできるようになります
  • user_idconfigurable から取得します。thread_id (Checkpointer のスコープ) とは独立にユーザーをスコープすることで、複数ユーザーのメモリが混ざりません
  • system_prompt でツールの使い分けを誘導します。これがないと LLM がツールを呼んでくれないことがあります

ステップ6: 動作確認 (セッション横断での記憶引き継ぎ)

ここから本記事の核心に入ります。新しいユーザー user-jiro で、セッションA → セッションB と切り替えながら、長期メモリが効いていることを確認します。

セッションA: 自己紹介で長期記憶に保存

SESSION_A = "session-A"
USER_ID = "user-jiro"

config_a = {
    "configurable": {
        "thread_id": SESSION_A,
        "user_id": USER_ID,
    }
}

print("=== セッションA: 自己紹介 ===")
result = agent.invoke(
    {"messages": [{"role": "user", "content": "私の名前は次郎です。山登りが趣味で、月に2回は奥多摩に行きます。"}]},
    config=config_a,
)
print(result["messages"][-1].content)

出力:

=== セッションA: 自己紹介 ===
奥多摩にはたくさんの魅力的な山がありますが、次郎さんはどちらの山によく登られるのでしょうか?
御岳山や大岳山、雲取山など、それぞれ違った魅力がありますよね。

エージェントは save_memory ツールを呼び、ユーザー情報を Store に保存しています(具体的に何を保存したかはステップ8で確認します)。

続けて別の話題を保存させます。

print("=== セッションA: 別の話題 ===")
result = agent.invoke(
    {"messages": [{"role": "user", "content": "好きな食べ物は寿司、特に光り物が好きです。"}]},
    config=config_a,
)
print(result["messages"][-1].content)

出力:

=== セッションA: 別の話題 ===
寿司がお好きなんですね!光り物は本当に美味しいですよね。
アジ、サバ、イワシ、コハダなど、それぞれ独特の風味と食感があって、寿司の醍醐味の一つだと思います。

山登りの後に美味しいお寿司を食べるのも格別でしょうね。
奥多摩から帰りに立ち寄る美味しいお寿司屋さんなどはありますか?

寿司・光り物の情報も保存されました。

セッションBに切り替えて記憶を引き継ぐか確認

ここが第2回の最大の見どころです。thread_id を変えて (= Checkpointer 的には完全に別セッション)、しかし user_id は同じにします。

SESSION_B = "session-B"

config_b = {
    "configurable": {
        "thread_id": SESSION_B,
        "user_id": USER_ID,  # user_idは同じ
    }
}

print("=== セッションB: 別スレッドで質問 ===")
result = agent.invoke(
    {"messages": [{"role": "user", "content": "週末におすすめの過ごし方を提案してください。"}]},
    config=config_b,
)
print(result["messages"][-1].content)

出力:

=== セッションB: 別スレッドで質問 ===
次郎さん、こんにちは!過去の会話を思い出しました。
山登りがお好きで、月に2回は奥多摩に行かれているんですよね。

そんな次郎さんにおすすめの週末の過ごし方をいくつか提案させていただきます:

## アウトドア系
1. **奥多摩での新しいルート開拓** - いつものコースとは違う山域を探索してみる
2. **高尾山周辺のハイキング** - 比較的軽めで、温泉も楽しめます
3. **丹沢方面への日帰り登山** - 奥多摩とは違った景色を楽しめます

## 山登り以外の選択肢
4. **登山用品店巡り** - 新しいギアをチェックしたり、次の山行の計画を立てる
5. **美味しい寿司屋探し** - 光り物の美味しいお店を新規開拓
6. **山の写真整理** - 過去の登山写真を整理してアルバム作り

(以下略)

完璧に動作しています。注目すべきは:

  • 「次郎さん、こんにちは!過去の会話を思い出しました」と明示的に過去の記憶を参照したことを表現している
  • 「山登りがお好きで、月に2回は奥多摩に行かれているんですよね」とセッションAで保存した情報を正確に引っ張ってきている
  • 提案リストの「5. 美味しい寿司屋探し - 光り物の美味しいお店を新規開拓」のように、複数の長期記憶を組み合わせて提案を構築している

thread_id が変わっているので Checkpointer の会話履歴は引き継がれていません。それなのに過去の情報が出てくるのは、recall_memories ツール経由で Store から思い出している証拠です。

ステップ7: 短期メモリと長期メモリの違いを動作で示す

セッションBの中で、続けてもう一つ質問してみます。

print("=== セッションB: 短期メモリの確認 (前ターンの直後の話題) ===")
result = agent.invoke(
    {"messages": [{"role": "user", "content": "さっき提案してくれた中で、3番目はなんでしたっけ?"}]},
    config=config_b,
)
print(result["messages"][-1].content)

出力:

=== セッションB: 短期メモリの確認 (前ターンの直後の話題) ===
3番目は「**丹沢方面への日帰り登山** - 奥多摩とは違った景色を楽しめます」でした。
丹沢は神奈川県にある山域で、奥多摩とはまた違った魅力がありますね。
(以下略)

同じ thread_id (session-B) なので、Checkpointer が直前の応答を覚えていて、「3番目は丹沢方面への日帰り登山」と正確に答えられています。これは 短期メモリ (会話履歴) が機能している例です。

次に、別のスレッドに切り替えて同じ質問をしてみます。

SESSION_C = "session-C"

config_c = {
    "configurable": {
        "thread_id": SESSION_C,  # 新しいthread
        "user_id": USER_ID,  # 同じユーザー
    }
}

print("=== セッションC: 別スレッドで「さっきの話」を聞く ===")
result = agent.invoke(
    {"messages": [{"role": "user", "content": "さっき提案してくれた中で、3番目はなんでしたっけ?"}]},
    config=config_c,
)
print(result["messages"][-1].content)

出力:

=== セッションC: 別スレッドで「さっきの話」を聞く ===
過去の記憶を確認しましたが、最近私がお出しした提案のリストについての記録が見つかりませんでした。

恐れ入りますが、どのような内容について提案をお求めでしたでしょうか?例えば:
- 山登りのスポットについて
- お食事の場所について
- その他の活動について

内容を教えていただければ、改めて提案させていただきます。

エージェントは recall_memories を試みた上で「提案のリストについての記録が見つかりませんでした」と正直に答えています。これは正しい挙動です。エージェントが直前にした提案リストは 会話の流れ (Checkpointer) に保存されたものであり、長期メモリ (Store) には保存されていないからです。一方で「山登りのスポットについて...」と選択肢を提示できているのは、Store には次郎さんが山登り好きだという長期記憶が残っているためです。

これで以下の対比が完璧に体現されました。

シナリオ thread_id user_id 結果 メカニズム
セッションA: 自己紹介 session-A user-jiro 山登り・奥多摩を保存 save_memory → Store
セッションA: 食べ物 session-A user-jiro 寿司・光り物を保存 save_memory → Store
セッションB: 提案 session-B user-jiro 過去の好みを思い出して提案 recall_memories → Store
セッションB: さっきの3番目 session-B user-jiro 「丹沢日帰り登山」と即答 Checkpointer (会話履歴)
セッションC: さっきの3番目 session-C user-jiro 知らないと回答 (でも長期記憶は使える) 別 thread → Checkpointer は使えない

Checkpointer (thread単位の会話履歴) と Store (user単位の長期記憶) が独立したレイヤーとして機能していることが完全に確認できました。

ステップ8: Lakebase内のStoreテーブルを覗く

何が Lakebase に保存されたのか、SQL で直接確認します。

import concurrent.futures

def inspect_store():
    with pool.connection() as conn:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT column_name, data_type 
                FROM information_schema.columns 
                WHERE table_name = 'store' AND table_schema = 'public'
                ORDER BY ordinal_position;
            """)
            store_schema = cur.fetchall()
            
            cur.execute("""
                SELECT prefix, key, value
                FROM store
                ORDER BY prefix, created_at;
            """)
            store_rows = cur.fetchall()
            
            cur.execute("SELECT COUNT(*) FROM store_vectors;")
            vec_count = cur.fetchone()[0]
            
            return store_schema, store_rows, vec_count

with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex:
    store_schema, store_rows, vec_count = ex.submit(inspect_store).result(timeout=15)

print("=== store table schema ===")
for col, dtype in store_schema:
    print(f"  {col}: {dtype}")

print(f"\n=== store contents ({len(store_rows)} rows) ===")
for prefix, key, value in store_rows:
    print(f"  prefix={prefix}")
    print(f"  key={key}")
    print(f"  value={value}")
    print()

print(f"=== store_vectors row count: {vec_count} ===")

出力:

=== store table schema ===
  prefix: text
  key: text
  value: jsonb
  created_at: timestamp with time zone
  updated_at: timestamp with time zone
  expires_at: timestamp with time zone
  ttl_minutes: integer

=== store contents (3 rows) ===
  prefix=memories.user-jiro
  key=memory-20260428-224738-262225
  value={'text': 'ユーザーの名前は次郎さん。趣味は山登りで、月に2回は奥多摩に行く。', 'saved_at': '2026-04-28T22:47:38.262250'}

  prefix=memories.user-jiro
  key=memory-20260428-224751-781576
  value={'text': '次郎さんの好きな食べ物は寿司、特に光り物(アジ、サバ、イワシなどの青魚)が好き。', 'saved_at': '2026-04-28T22:47:51.781599'}

  prefix=memories.user-taro
  key=memory-001
  value={'text': '私は温泉が好きで、特に箱根によく行きます'}

=== store_vectors row count: 3 ===

いくつか観察できることがあります。

prefixmemories.user-jiro / memories.user-taro という形でユーザー単位のスコープが効いています。これによって user-jiro の検索結果に user-taro のメモリが混ざることはありません。

value は JSONB なので、text だけでなく任意のメタデータ (本記事では saved_at を一緒に保存) を柔軟に格納できます。

エージェントが保存したテキストを見ると、ユーザーが「私の名前は次郎です」と言ったのを「ユーザーの名前は次郎さん」と三人称化して保存しています。LLMが「将来読み返す自分」を意識して整形しているのが見て取れて興味深いです。

expires_atttl_minutes のフィールドがあるので、メモリに有効期限を設定して自動失効させる運用も可能です。本記事では使っていませんが、長期間のセッションで「古い情報を自動的に忘れる」運用が必要な場合に便利です。

store_vectors テーブルには対応する3件のEmbeddingベクトル (1024次元) が保存されており、これが recall_memories での意味検索のインデックスとして機能しています。

これらのデータはLakebase GUIからも確認できます。

Screenshot 2026-04-29 at 7.57.15.png

ハマりどころのまとめ

第1回からの教訓も含めて、本記事で踏んだハマりどころです。

databricks-langchain を入れると langgraph がダウングレードされて ImportError: cannot import name 'ExecutionInfo' が出る (第1回からの継続課題)。これは --upgradelanggraph を再インストールすれば解消する。

Qwen3-Embedding が稀に int 値を含むベクトルを返す。psycopg + pgvector に渡す前に明示的に float() でキャストしないと「cannot dump lists of mixed types」エラーになる。Embedding API の出力を信用せず、ベクトル DB に投入する前に正規化する習慣をつけておくと良い。

PostgresStore.search() の戻り値の score は cosine similarity ではなく 距離指標 として解釈すべき。値が高ければ類似度が高い、と直感に頼らず、いくつかのサンプルでスコアの分布を確認してからしきい値を決めるのが安全。

save_memory ツールが LLM によって自動で呼ばれるかは system prompt の書き方に強く依存 する。ツールのdocstringだけでは LLM が呼んでくれないことが多いので、prompt で明示的に「ユーザーが自分について話したら save_memory を使う」と指示するのが効果的。

store テーブルの prefix["memories", "user-taro"] のようなタプルが内部で memories.user-taro という文字列にエンコードされて保存される。SQLで直接覗く時は文字列形式で見える点に注意。

OAuthトークンの60分失効は第1回と同じく注意が必要。本番運用ではConnectionPoolconnection_class をカスタマイズしてトークン自動更新を実装するのが定石。

次回予告

連載第3回では、評価ハーネス層を扱います。具体的には:

  • MLflow Tracing でエージェントの全ステップを記録
  • 過去のトレースから評価データセットを自動生成
  • カスタムジャッジ (make_judge) で「メモリを適切に活用したか」を評価
  • 継続的改善ループ (悪いケースをデータセットに追加 → 改善 → 再評価)

メモリだけ持っていても「ちゃんと活用できているか」を測れないと改善できません。ステートフルエージェントの観測性と評価を、Databricks のフルスタックで実装する回になります。

本記事の構成を Unity Catalog 登録 + デプロイ + Agent Evaluation までフルパスで実装する場合は、 Databricks Lakebaseを用いたステートフルAIエージェント を参照してください。ResponsesAgent でラップする部分や agents.deploy() の流れがそのまま参考になります。

まとめ

  • LangGraph の PostgresStore を使うと、thread_id を跨いで永続化されるユーザー単位の長期メモリが実装できる
  • Checkpointer (短期、会話履歴) と Store (長期、ユーザー知識) は独立したレイヤーで、両方を併用するエージェントが構成できる
  • pgvector + Qwen3-Embedding-0.6B で日本語の意味検索が高品質に動作する
  • save_memory / recall_memories ツール + system prompt で誘導すれば、LLM が自然に長期記憶の保存・参照を判断する
  • Lakebase Autoscaling は前回作ったプロジェクトの上に拡張的にテーブルを増やすだけで運用でき、データの永続性も体感できる

参考リンク

はじめてのDatabricks

はじめてのDatabricks

Databricks無料トライアル

Databricks無料トライアル

4
2
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
4
2

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?