すべてのプロダクト
Search
ドキュメントセンター

Data Lake Formation:PyPaimon による DLF Paimon テーブルに対するスタンドアロン SQL クエリの実行

最終更新日:Jun 12, 2026

PyPaimon に組み込まれた SQL エンジンを使用すると、ローカルの Python 環境から直接、Data Lake Formation (DLF) のカタログで管理されている Paimon テーブルにクエリを実行できます。Flink や Spark クラスターは不要なため、このアプローチはデータ探索、モデルのデバッグ、レポート生成などのアドホックタスクに適しています。

前提条件

依存関係のインストール

  1. PyPaimon オフラインパッケージをダウンロードします: pypaimon-1.5.dev20260608.tar.gz

  2. PyPaimon パッケージとその他の依存関係をインストールします:

    pip install \
        pypaimon-1.5.dev20260608.tar.gz \  # PyPaimon オフラインパッケージ
        datafusion \                        # SQL エンジン
        pypaimon-rust \                     # DataFusion 統合
        pyarrow pandas requests             # データ形式とネットワーク
  3. 次のコードを実行して、datafusion と PaimonCatalog が利用可能であることを確認します:

    from datafusion import SessionContext
    from pypaimon_rust.datafusion import PaimonCatalog
    import datafusion
    
    print("datafusion:", datafusion.__version__)
    print("PaimonCatalog OK:", PaimonCatalog)

カタログ接続の設定

コードで DLF カタログの接続パラメーターを定義します:

CATALOG_OPTIONS = {
    "metastore": "rest",
    "uri": "http://<DLF-ENDPOINT>",
    "warehouse": "<YOUR-CATALOG>",
    "token.provider": "dlf",
    "dlf.region": "<REGION-ID>",
    "dlf.access-key-id": "<ACCESS-KEY-ID>",
    "dlf.access-key-secret": "<ACCESS-KEY-SECRET>",
    "dlf.oss-endpoint": "<OSS-ENDPOINT>",
}

次の表で各パラメーターについて説明します。

パラメーター

説明

metastore

固定値: rest

uri

DLF Paimon REST エンドポイント。詳細については、「Endpoints」を参照してください。VPC アクセスは HTTP と HTTPS の両方をサポートします。パブリックアクセスには HTTPS が必要です。

warehouse

DLF カタログ名。

token.provider

固定値: AccessKey 認証を使用する場合、dlf です。

dlf.region

DLF がデプロイされているリージョン ID (例: cn-hangzhou)。詳細については、「エンドポイント」をご参照ください。

dlf.access-key-id / dlf.access-key-secret

お使いの Alibaba Cloud アカウントまたは RAM ユーザーの AccessKey ペア。

dlf.oss-endpoint

(オプション) OSS パブリックエンドポイント。パブリックアクセスに必要です。VPC アクセスの場合、省略できます。

権限要件

RAM ユーザーとして DLF にアクセスする場合、ユーザーが必要な API 権限とデータ権限を持っていることを確認してください。詳細については、「権限の設定」を参照してください。

DataFusion への Paimon カタログの登録

PaimonCatalog は DataFusion カタログプロバイダープロトコルを実装しています。SessionContext に登録すると、3部構成の名前 paimon.<database>.<table> を使用して Paimon テーブルを参照できます。

from datafusion import SessionContext
from pypaimon_rust.datafusion import PaimonCatalog

ctx = SessionContext()
paimon = PaimonCatalog(CATALOG_OPTIONS)
ctx.register_catalog_provider("paimon", paimon)

print("catalogs:", ctx.catalog_names())
# -> {'paimon', 'datafusion'}

# paimon カタログ内のすべてのデータベース (スキーマ) を一覧表示
print("databases:", ctx.catalog("paimon").names())

# default データベース内のすべてのテーブルを一覧表示
print("tables in default:", ctx.catalog("paimon").schema("default").names())

Paimon テーブルへのクエリ

以下の例では、employees テーブル (列 idnamedeptsalary) と departments テーブルがすでに存在することを前提としています。テーブル作成とデータ取り込みのコードについては、このドキュメントの最後にある [完全な例](#complete-example) をご参照ください。

クエリ結果は、.to_pylist().to_arrow_table().to_pandas() などのメソッドを使用して、さまざまなデータ形式に変換できます。

SELECT クエリ

rows = ctx.sql("SELECT * FROM paimon.default.employees ORDER BY id").to_pylist()
print(rows[0])
# -> {'id': 1, 'name': 'Alice', 'dept': 'eng', 'salary': 9000}

WHERE フィルター、カラムプロジェクション、および ORDER BY

eng = ctx.sql("""
    SELECT id, name, salary
    FROM paimon.default.employees
    WHERE dept = 'eng'
    ORDER BY salary DESC
""").to_pylist()
print(eng)
# -> [{'id': 4, 'name': 'Dan', 'salary': 9500},
#     {'id': 1, 'name': 'Alice', 'salary': 9000},
#     {'id': 2, 'name': 'Bob', 'salary': 8500}]

GROUP BY 集約

agg = ctx.sql("""
    SELECT dept, COUNT(*) AS n, AVG(salary) AS avg_sal
    FROM paimon.default.employees
    GROUP BY dept
    ORDER BY dept
""").to_pylist()
print(agg)
# -> [{'dept': 'eng',   'n': 3, 'avg_sal': 9000.0},
#     {'dept': 'ops',   'n': 1, 'avg_sal': 6500.0},
#     {'dept': 'sales', 'n': 2, 'avg_sal': 7100.0}]

JOIN

自己結合の例:

result = ctx.sql("""
    SELECT a.id AS id_a, b.id AS id_b, a.name
    FROM paimon.default.employees a
    JOIN paimon.default.employees b
      ON a.id = b.id
    WHERE a.dept = 'eng'
    ORDER BY a.id
""").to_pylist()

クロステーブル結合の例:

result = ctx.sql("""
    SELECT e.id, e.name, e.dept, d.dept_full
    FROM paimon.default.employees e
    JOIN paimon.default.departments d
      ON e.dept = d.dept
    WHERE e.dept = 'eng'
    ORDER BY e.id
""").to_pylist()

クエリ結果の変換

クエリ結果を複数の形式に変換して、さまざまな下流の処理フレームワークと統合できます。

Arrow テーブルへの変換

arrow_tbl = ctx.sql(
    "SELECT id, salary FROM paimon.default.employees ORDER BY id"
).to_arrow_table()
print(arrow_tbl.num_rows, "rows")
print(arrow_tbl.schema)

Pandas DataFrame への変換

df = ctx.sql(
    "SELECT * FROM paimon.default.employees ORDER BY id"
).to_pandas()
print(df.head())

Daft との統合

Daft フレームワークを使用している場合は、Arrow テーブルを Daft DataFrame に変換してさらに処理を進めます。Daft と DLF Paimon テーブルの直接統合については、「Daft を DLF Paimon テーブルで使用する」を参照してください。

import daft

daft_df = daft.from_arrow(arrow_tbl)
daft_df.show()

完全な例

次の例は、PyPaimon を使用してテーブルを作成し、データを書き込み、次に DataFusion にカタログを登録して SQL クエリを実行するという、エンドツーエンドのワークフローを示しています。

from __future__ import annotations

import os
from datetime import datetime

import pyarrow as pa
from datafusion import SessionContext
from pypaimon_rust.datafusion import PaimonCatalog
from pypaimon import CatalogFactory, Schema


CATALOG_OPTIONS = {
    "metastore": "rest",
    "uri": "http://<DLF-ENDPOINT>",
    "warehouse": "<YOUR-CATALOG>",
    "token.provider": "dlf",
    "dlf.region": "<REGION-ID>",
    "dlf.access-key-id": os.environ["DLF_ACCESS_KEY_ID"],
    "dlf.access-key-secret": os.environ["DLF_ACCESS_KEY_SECRET"],
    "dlf.oss-endpoint": "<OSS-ENDPOINT>",
}

DATABASE = "default"


def main() -> None:
    # 1) PyPaimon を使用してテーブルを作成し、データを書き込む
    pcat = CatalogFactory.create(CATALOG_OPTIONS)
    pcat.create_database(DATABASE, True)

    table_name = f"sql_demo_{datetime.now().strftime('%Y%m%d_%H%M%S')}"
    full_name = f"{DATABASE}.{table_name}"
    pa_schema = pa.schema([
        pa.field("id",     pa.int32(), nullable=False),
        pa.field("name",   pa.string()),
        pa.field("dept",   pa.string()),
        pa.field("salary", pa.int64()),
    ])
    schema = Schema.from_pyarrow_schema(
        pa_schema,
        options={"bucket": "-1", "file.format": "parquet"},
    )
    pcat.create_table(full_name, schema, True)

    tbl = pcat.get_table(full_name)
    wb = tbl.new_batch_write_builder()
    w = wb.new_write()
    w.write_arrow(pa.Table.from_pydict({
        "id":     [1, 2, 3, 4, 5, 6],
        "name":   ["Alice", "Bob", "Carol", "Dan", "Eve", "Frank"],
        "dept":   ["eng", "eng", "sales", "eng", "sales", "ops"],
        "salary": [9000, 8500, 7000, 9500, 7200, 6500],
    }, schema=pa_schema))
    wb.new_commit().commit(w.prepare_commit())
    w.close()
    print(f"created + wrote 6 rows to {full_name}")

    # 2) DataFusion にカタログを登録する
    ctx = SessionContext()
    ctx.register_catalog_provider("paimon", PaimonCatalog(CATALOG_OPTIONS))

    # 3) SQL クエリを実行する
    rows = ctx.sql(
        f"SELECT * FROM paimon.{DATABASE}.{table_name} ORDER BY id"
    ).to_pylist()
    print("all rows:", rows)

    eng = ctx.sql(
        f"SELECT id, name, salary FROM paimon.{DATABASE}.{table_name} "
        f"WHERE dept = 'eng' ORDER BY salary DESC"
    ).to_pylist()
    print("eng:", eng)

    agg = ctx.sql(
        f"SELECT dept, COUNT(*) AS n, AVG(salary) AS avg_sal "
        f"FROM paimon.{DATABASE}.{table_name} GROUP BY dept ORDER BY dept"
    ).to_pylist()
    print("by dept:", agg)


if __name__ == "__main__":
    main()

PyPaimon の SQL CLI の使用

PyPaimon には、組み込みの SQL CLI が含まれています。YAML 設定ファイルでカタログ接続を指定した後、対話型 REPL を開始するか、単一の SQL ステートメントを実行できます。

カタログ設定を catalog.yaml として保存します:

metastore: rest
uri: http://<DLF-VPC-ENDPOINT>
warehouse: <YOUR-CATALOG>
token.provider: dlf
dlf.region: <REGION-ID>
dlf.access-key-id: <ACCESS-KEY-ID>
dlf.access-key-secret: <ACCESS-KEY-SECRET>
dlf.oss-endpoint: <OSS-ENDPOINT>
説明

セキュリティに関する推奨事項: AccessKey 認証情報は、設定ファイルにプレーンテキストで格納するのではなく、環境変数を介して注入してください。 たとえば、export DLF_ACCESS_KEY_ID=xxxexport DLF_ACCESS_KEY_SECRET=xxx を設定し、envsubst で YAML テンプレートをレンダリングします。 本番環境では、STS 一時的認証情報を使用してください。

paimon コマンドを使用してクエリを実行します:

# 対話型 REPL を開始
paimon --config /path/to/catalog.yaml sql

# 単一の SQL ステートメントを実行
paimon --config /path/to/catalog.yaml sql "SELECT * FROM default.employees LIMIT 5"

# 出力形式を指定 (デフォルト: table、json もサポート)
paimon --config /path/to/catalog.yaml sql --format json "SELECT * FROM default.employees"

出力形式や結果のエクスポートなどの追加オプションについては、paimon sql --help を実行すると、すべてのパラメーターリストが表示されます。

SQL クエリインターフェイスの比較

PyPaimon は、使用パターンが異なる 2 つの SQL クエリインターフェイスを提供しています。このドキュメントでは、PaimonCatalog + DataFusion 方式について説明します。SQLContext 方式については、PyPaimon API ドキュメントをご参照ください。

インターフェイス

テーブル参照

戻り値の型

ユースケース

PaimonCatalog + datafusion.SessionContext (本書で説明)

3部構成の名前: paimon.<database>.<table>

.to_pylist().to_arrow_table().to_pandas() などの変換メソッドを持つ DataFusion データフレーム

DataFusion エコシステムとの統合、または結果を Arrow、Pandas、Polars に直接変換する場合

pypaimon.SQLContext

set_current_catalog または set_current_database を呼び出した後の単一パートのテーブル名

list[pyarrow.RecordBatch] (Pandas/Polars への手動変換が必要です)

SQL を実行して RecordBatch の結果を取得するだけでよい軽量なシナリオ