PyPaimon に組み込まれた SQL エンジンを使用すると、ローカルの Python 環境から直接、Data Lake Formation (DLF) のカタログで管理されている Paimon テーブルにクエリを実行できます。Flink や Spark クラスターは不要なため、このアプローチはデータ探索、モデルのデバッグ、レポート生成などのアドホックタスクに適しています。
前提条件
依存関係のインストール
PyPaimon オフラインパッケージをダウンロードします: pypaimon-1.5.dev20260608.tar.gz
PyPaimon パッケージとその他の依存関係をインストールします:
pip install \ pypaimon-1.5.dev20260608.tar.gz \ # PyPaimon オフラインパッケージ datafusion \ # SQL エンジン pypaimon-rust \ # DataFusion 統合 pyarrow pandas requests # データ形式とネットワーク次のコードを実行して、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 | 固定値: |
uri | DLF Paimon REST エンドポイント。詳細については、「Endpoints」を参照してください。VPC アクセスは HTTP と HTTPS の両方をサポートします。パブリックアクセスには HTTPS が必要です。 |
warehouse | DLF カタログ名。 |
token.provider | 固定値: AccessKey 認証を使用する場合、 |
dlf.region | DLF がデプロイされているリージョン ID (例: |
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 テーブル (列 id、name、dept、salary) と 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=xxx と export 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 ドキュメントをご参照ください。
インターフェイス | テーブル参照 | 戻り値の型 | ユースケース |
| 3部構成の名前: |
| DataFusion エコシステムとの統合、または結果を Arrow、Pandas、Polars に直接変換する場合 |
|
|
| SQL を実行して RecordBatch の結果を取得するだけでよい軽量なシナリオ |