0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Databricks Lakebase Searchを使ったRAGエージェントをFree Edition上で作ってみる

0
Posted at

はじめに

Databricks Lakebaseにて、Lakebase Searchという検索機能が提供されています。

Lakebase Search(lakebase_vectorとlakebase_text拡張)はベータ版です(2026年7月時点)。

Lakebase SearchはPostgreSQLベースのLakebaseにlakebase_vector拡張を入れることで、ANNインデックス(pgvectorの進化版)とBM25ハイブリッド検索が使えるサービスです。

Databricksは既にAI Searchが提供されているのですが、エンドポイントの常時起動が必要など、コスト面で気になる点もあります。
Lakebaseはゼロ・スケール可能ということもあり、こちらでいわゆるAgentic RAGが利用できるかを試してみました。

今回は以下の流れで、サンプルデータを使ってレビュー検索RAGエージェントを作ってみます。

  1. レビューデータに埋め込みを入れたパイプラインを作る
  2. Lakebaseプロジェクトをセットアップし、Lakehouse Syncでデータを同期
  3. Lakebase Searchのインデックスを作成
  4. LangChainでRAGエージェントを組み立てる

Lakebase Searchとは

DatabricksのLakebaseプロジェクトから利用できるPostgreSQL拡張機能で、ベクトル検索とハイブリッド検索を提供します。

主な特徴は以下の通りです:

  • lakebase_vector拡張: pgvectorベースで、ANN(近似最近傍探索)インデックスをサポート
  • lakebase_text拡張: BM25インデックスによる全文検索
  • ハイブリッド検索: ベクトル検索とキーワード検索をRRF(Reciprocal Rank Fusion)で統合

また、以下のblogで詳しく解説されています。

Lakebaseはインスタンスをゼロスケールできるので、コストパフォーマンス高く柔軟にベクトル検索を行える仕組として優れているのではないかと思います。

検証環境

  • Databricks Free Editionを使います。
  • ノートブックのサーバレス環境はStandard v5を利用しました。

今回作るもの

レビューコメントをベクトル化し、自然言語クエリで関連レビューを検索できるRAGエージェントを作ります。

処理フローは以下のような流れです。

[Unity Catalog: samples.wanderbricks.reviews]
  ↓ (Spark Declarative Pipeline)
[埋め込みベクトル付きMaterialized View]
  ↓ (Lakehouse Sync)
[Lakebase: reviews_with_embeddings_synced]
  ↓ (ETL + Index作成)
[検索テーブル: reviews_search + ANNインデックス]
  ↓
[LangChain AgentでRAG]

Step1. 埋め込みパイプラインの作成

まず、レビューデータに埋め込みベクトルを追加するパイプラインを作成します。

パイプライン構成

Spark宣言型パイプラインを作成し、transformations/bronze/reviews.pyを作成。

transformations/bronze/reviews.py
from pyspark import pipelines as dp
from pyspark.sql.functions import expr, col

@dp.materialized_view(
    comment="Bronze層: samples.wanderbricks.reviewsからの生レビューデータ"
)
def bronze_reviews():
    return spark.read.table("samples.wanderbricks.reviews")


@dp.materialized_view(
    comment="埋め込みベクトル付きレビュー - 先頭5件のコメントをdatabricks-qwen3-embedding-0-6bでベクトル化"
)
def reviews_with_embeddings():
    return (
        spark.read.table("bronze_reviews")
        .filter(col("comment").isNotNull())
        .limit(1000) # データ量節約のため、1000件のみ対象
        .withColumn(
            "comment_embedding",
            expr("ai_query('databricks-qwen3-embedding-0-6b', comment)")
        )
    )

今回はサンプルデータであるsamples.wanderbricks.reviewsを利用しています。
これは架空のユーザレビューを含んだテーブルです。

パイプライン内ではai_query関数で埋め込みモデルを使ってベクトルデータを追加したMaterialized View(MV)を作成するよう定義しています。
このパイプラインを実行して、workspace.lakebase_search_testにMVを作成します。

Step2. Lakebaseプロジェクト作成とLakehouse Sync

次に、Lakebaseプロジェクトを作成し、Unity CatalogテーブルをLakebaseに同期します。

2-1. プロジェクト作成

ノートブックを作成し、まずは以下を実行してLakebaseのプロジェクトを作成します。

from databricks.sdk import WorkspaceClient
from databricks.sdk.service.postgres import Project, ProjectSpec

w = WorkspaceClient()

project_id = "lakebase-search-test"

operation = w.postgres.create_project(
    project=Project(
        spec=ProjectSpec(
            display_name=project_id,
            pg_version=17,
        )
    ),
    project_id=project_id
)

result = operation.wait()
print(f"Created project: {result.name}")

2-2. Synced Tableの作成

Lakehouse SyncでUnity Catalog上のMaterialized ViewをLakebase側に同期します:

from databricks.sdk.service.postgres import (
    SyncedTable,
    SyncedTableSyncedTableSpec,
    SyncedTableSyncedTableSpecSyncedTableSchedulingPolicy,
)

table_full_name = "workspace.lakebase_search_test.reviews_with_embeddings"
branch = f"projects/{project_id}/branches/production"
synced_table_id = "workspace.lakebase_search_test.reviews_with_embeddings_synced"

synced_table = w.postgres.create_synced_table(
    synced_table=SyncedTable(
        spec=SyncedTableSyncedTableSpec(
            source_table_full_name=table_full_name,
            branch=branch,
            primary_key_columns=["booking_id"],
            scheduling_policy=SyncedTableSyncedTableSpecSyncedTableSchedulingPolicy.SNAPSHOT,
            postgres_database="postgres_database",
            create_database_objects_if_missing=True,
        )
    ),
    synced_table_id=synced_table_id,
).wait()

print(f"Synced table created: {synced_table.name}")

scheduling_policy=SNAPSHOTで定期的に全件同期します。増分同期(CDC)も選べますが、今回はシンプルにスナップショットでいきます。

これでLakehouseのMVがLakebase側のテーブルとして同期できました。

Step3. Lakebase Searchのセットアップとインデックス作成

Lakebase側でテーブルとインデックスを作成します。
以下のSQLをLakebaseのSQL Editor上で実行していきます。

3-1. 拡張のインストール

まず、Lakebaseプロジェクトで拡張を有効化します。

Lakebaseのプロジェクト設定画面を開き、Lakebase Searchセクションの「Enable Lakebase Search」ボタンを押して、Lakebase Searchを有効化します。

image.png

次にSQL Editorで拡張機能をインストールします。

CREATE EXTENSION IF NOT EXISTS lakebase_vector CASCADE;
CREATE EXTENSION IF NOT EXISTS lakebase_text;

3-2. 検索用テーブルの作成

LakehouseからSyncしたテーブルはそのままではインデックスを作成できないため、別にインデックス用のテーブルを作成し、データをコピーします。

全データのコピーは正直イケてないように思います。
LTAPが来ればこの問題も解消するのではないかと思うのですが、いいやり方が知りたいですね。

CREATE TABLE IF NOT EXISTS lakebase_search_test.reviews_search (
  review_id    BIGINT,
  booking_id   BIGINT,
  property_id  BIGINT,
  user_id      BIGINT,
  rating       DOUBLE PRECISION,
  comment      TEXT,
  created_at   TIMESTAMPTZ,
  updated_at   TIMESTAMPTZ,
  is_deleted   BOOLEAN,
  comment_embedding VECTOR(1024),
  comment_tsv  TSVECTOR
);

INSERT INTO lakebase_search_test.reviews_search
SELECT
  review_id, booking_id, property_id, user_id, rating,
  comment, created_at, updated_at, is_deleted,
  (comment_embedding::text)::vector            AS comment_embedding,
  to_tsvector('english', coalesce(comment,'')) AS comment_tsv
FROM lakebase_search_test.reviews_with_embeddings_synced
WHERE is_deleted = false;

3-3. ANNインデックスとBM25インデックス

ベクトル検索用のインデックスとキーワード検索用のBM25インデックスを作成します。

-- ANNインデックス(コサイン距離)
CREATE INDEX ON lakebase_search_test.reviews_search 
  USING lakebase_ann (comment_embedding vector_cosine_ops);

-- BM25インデックス
CREATE INDEX reviews_comment_bm25 ON lakebase_search_test.reviews_search 
  USING lakebase_bm25 (comment_tsv);

埋め込みモデルはコサイン距離での運用が標準なのでvector_cosine_opsを指定しています。

Step4. RAGエージェントの作成

検索用テーブルができたので最後に、LangChainでRAGエージェントを組み立てます。
今回は試験用にノートブックを作成して以下の処理を実行していきます。

4-1. ライブラリインストール

langchainなど、必要なパッケージをインストール。

%uv pip install langchain>=1.3.11 databricks-langchain>=0.20.0 psycopg[pool]>=3.3.4 pgvector>=0.4.2 databricks-sdk>=0.120.0 mlflow>=3.14.0
%restart_python

4-2. Lakebase接続プール設定

接続の度にトークンを再生成するカスタム接続クラスを使ってlakebase用の接続プールを作成します。
参考にする場合は、ユーザ名やホスト名は環境に合わせて変更ください。

import os
import psycopg
from psycopg_pool import ConnectionPool
from psycopg.rows import dict_row
from databricks.sdk import WorkspaceClient

w = WorkspaceClient()

ENDPOINT_NAME = "projects/lakebase-search-test/branches/production/endpoints/primary"

class CustomConnection(psycopg.Connection):
    @classmethod
    def connect(cls, conninfo="", **kwargs):
        endpoint = ENDPOINT_NAME
        credential = w.postgres.generate_database_credential(endpoint=endpoint)
        kwargs["password"] = credential.token
        return super().connect(conninfo, **kwargs)

username = "xxx@xxxx.com" # ユーザ名
host = "ep-xxxxxxxxxx.us-east-2.cloud.databricks.com" # ホストアドレス
port = 5432
database = "postgres_database"

pool = ConnectionPool(
    conninfo=f"dbname={database} user={username} host={host} sslmode=require",
    connection_class=CustomConnection,
    min_size=1,
    max_size=10,
    open=True,
    kwargs={"row_factory": dict_row},
)

4-3. 検索ツールの実装

Lakebase Searchを使うツールを定義します。
実体は、埋め込みモデルを呼び出してベクトルデータを取り出す関数と、Lakebase Searchを使うSQLを実行するツールで構成されています。
今回はキーワード検索を使わず、ベクトル検索のみを実行するものとしました。

キーワード検索を含めたハイブリッド検索を試す予定でしたが、うまくいかず。原因確認中です。

from databricks_langchain import ChatDatabricks, DatabricksEmbeddings
from langchain_core.tools import tool
from langchain.agents import create_agent
import json

EMBED_ENDPOINT = "databricks-qwen3-embedding-0-6b"
LLM_ENDPOINT = "databricks-qwen3-next-80b-a3b-instruct"

embeddings = DatabricksEmbeddings(endpoint=EMBED_ENDPOINT)

def embed_query(text: str) -> list[float]:
    prompt = f"Instruct: Given a search query, retrieve relevant review comments\\nQuery: {text}"
    return embeddings.embed_query(prompt)

_HYBRID_SQL = """
WITH vec AS (
    SELECT review_id, row_number() OVER (ORDER BY dist) AS rank
    FROM (
        SELECT review_id, comment_embedding <=> %(qvec)s AS dist
        FROM lakebase_search_test.reviews_search
        ORDER BY dist
        LIMIT %(cand)s
    ) v
)
SELECT r.review_id, r.property_id, r.rating, r.comment, vec.rank
FROM lakebase_search_test.reviews_search r JOIN vec ON vec.review_id = r.review_id
ORDER BY rank
LIMIT %(limit)s;
"""

@tool
def search_reviews(query: str, limit: int = 10) -> str:
    """Search review comments by meaning.

    Use this whenever the user asks about what guests said, sentiment, or specifics
    mentioned in reviews. Returns matching reviews as a JSON array.

    Args:
        query: Natural-language search query in English.
        limit: Number of reviews to return  (default 10).
    """
    params = {
        "qvec": str(embed_query(query)),
        "qtext": query,
        "cand": 50,
        "limit": limit,
    }
    with pool.connection() as conn, conn.cursor() as cursor:
        cursor.execute(_HYBRID_SQL, params)
        rows = cursor.fetchall()
    return json.dumps([dict(r) for r in rows], ensure_ascii=False, default=str)

4-4. エージェントの組み立て

最後に、langchainを使ってエージェントを作成します。
ツールとして先ほど定義した検索ツールを利用します。

llm = ChatDatabricks(endpoint=LLM_ENDPOINT, temperature=0)

agent = create_agent(
    model=llm,
    tools=[search_reviews],
    system_prompt=(
        "You are a review-analysis assistant. When the user asks about guest feedback, "
        "call search_reviews to retrieve relevant reviews, then answer grounded in them. "
        "Cite review_id values you relied on."
    ),
)

使ってみる

実際にエージェントを動かしてみます。

import mlflow
mlflow.langchain.autolog()

result = agent.invoke({"messages": [
    {"role": "user", "content": "ネガティブな意見にはどんなものがあるか教えて"}
]})
print(result["messages"][-1].content)
実行結果
ネガティブな意見として、複数のレビューで共通して挙げられているのは「**清潔さに欠けていた**」という点です。以下のように、多くのゲストが「Disappointing experience. The place wasn’t clean.(がっかりした。場所が清潔ではなかった。)」と明確に不満を述べています。

この意見は、レビューID 285、243、568、640、421、492、704、179、944、333 など、10件以上のレビューに繰り返し登場しており、清潔さの問題が主要なネガティブフィードバックとして顕著です。

その他の具体的な不満については、現在のデータではこの1点に集中しています。

元データは同一コメントが繰り返し使われているようで、idは別物ですがコメント内容は同じものが引っかかってきました。
とはいえ、検索自体は正常にできていそうです。

MLflowのトレースを見ても正常にツールが呼ばれて検索できていました。

image.png

まとめ

Lakebase SearchでベクトルRAGエージェントを試してみました。

Lakehouseで構築したベクトル情報を含むテーブルを一定利用してLakebase上にインデックスrを作れるのは魅力的だと思いました。AI Search用のエンドポイントを用意しなくていいのはメリットです。

とはいえ、Lakehouse Syncの同期方法検討やLakebase上でデータコピーが必要など、実運用ではまだ検証すべき点が残っています。
なんとなくですが、AI SearchとLakebase Searchはマージしていくのではないかと考えていますので、この先の発展が楽しみです。(DAIS2026でそういう発表を期待していたのですが)

Databricks Apps上でMCPサーバ化すればさらに使い勝手も上がるように思うので、このあたりも試してみたいと思います。

0
0
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
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?