Daft は、pypaimon の read_paimon / write_paimon 関数を介して、Data Lake Formation (DLF) カタログで管理されている Paimon テーブルの読み取りと書き込みを行います。述語プッシュダウン、プロジェクションプッシュダウン、タイムトラベルクエリをサポートします。このトピックでは、最小限のエンドツーエンドの例を紹介します。
BLOB / list[binary] マルチモーダル列を操作する場合、またはマルチモーダル ETL を実行する場合は、Daft を使用した DLF Paimon でのマルチモーダル自動運転データの管理 をご参照ください。
前提条件
ステップ 1:依存関係のインストール
-
pypaimon オフラインインストールパッケージのダウンロード:pypaimon-1.5.dev20260608.tar.gz
-
Python 依存関係のインストール
オフラインパッケージが含まれるディレクトリで、次のコマンドを実行してすべての依存関係をインストールします。
pip install daft pyarrow pandas requests pypaimon-1.5.dev20260608.tar.gz -
環境の確認
次のコードを実行して、環境が正しく設定されていることを確認します。
import daft from pypaimon.daft import read_paimon, write_paimon print("daft:", daft.__version__, "pypaimon.daft: OK")
ステップ 2:カタログ接続の設定
カタログ接続パラメータを次のように設定します。
CATALOG_OPTIONS = {
"metastore": "rest",
"uri": "http://<DLF-ENDPOINT>", // パブリックインターネット経由では HTTPS が必要です
"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>", // VPC ではオプション、パブリックインターネットでは必須です
}
|
パラメータ |
説明 |
|
metastore |
|
|
uri |
DLF Paimon REST エンドポイント。詳細については、「エンドポイント」をご参照ください。VPC アクセスでは HTTP と HTTPS の両方をサポートしますが、パブリックインターネットアクセスでは HTTPS が必要です。VPC アクセスの例: http://ap-southeast-1-vpc.dlf.aliyuncs.com |
|
warehouse |
DLF カタログの名前。 |
|
token.provider |
AccessKey 認証の場合は |
|
dlf.region |
DLF リージョン ID (例: |
|
dlf.access-key-id / dlf.access-key-secret |
DLF 認証用のアクセスキーペア。 |
|
dlf.security-token |
(オプション) STS セキュリティトークン。 |
|
dlf.oss-endpoint |
(オプション) OSS パブリックエンドポイント。 |
データの書き込み
新しいテーブルの作成
Daft でテーブルの読み取りまたは書き込みを行う前に、PyPaimon の create_table() を使用してテーブルを作成します。
from datetime import datetime
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
catalog = CatalogFactory.create(CATALOG_OPTIONS)
catalog.create_database("default", True)
table_name = "default.test_daft_paimon_" + datetime.now().strftime("%Y%m%d_%H%M%S")
pa_schema = pa.schema([
pa.field("id", pa.int32(), nullable=False),
pa.field("name", pa.string()),
pa.field("score", pa.int64()),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
options={"bucket": "-1", "file.format": "parquet"},
)
catalog.create_table(table_name, schema, True)
テーブルへの書き込み
write_paimon を使用して、Daft データフレームを既存の Paimon テーブルに書き込みます。
import daft
from pypaimon.daft import write_paimon
append_df = daft.from_pydict({
"id": [4, 5],
"name": ["Dan", "Eve"],
"score": [77, 92],
})
write_paimon(append_df, table_name, CATALOG_OPTIONS, mode="append")
mode パラメーターは 2 つの値を受け入れます:
-
"append":テーブルに行を追加します。 -
"overwrite":テーブル内のすべての既存データを上書きします。
既存テーブルの読み取り
基本的な読み取り
read_paimon を呼び出して、テーブルを Daft データフレームに読み込みます。
from pypaimon.daft import read_paimon
df = read_paimon(table_name, CATALOG_OPTIONS)
df.show()
述語プッシュダウンとプロジェクションプッシュダウン
Daft は、フィルター述語と列プロジェクションを Paimon スキャンにプッシュダウンします。
import daft
df = (
read_paimon(table_name, CATALOG_OPTIONS)
.where(daft.col("score") >= 88)
.select("id", "name")
)
df.show()
タイムトラベル
履歴バージョンを読み取るには、snapshot_id または tag_name を渡します。
# スナップショット ID による指定
df = read_paimon(table_name, CATALOG_OPTIONS, snapshot_id=42)
# タグ名による指定
df = read_paimon(table_name, CATALOG_OPTIONS, tag_name="v1")
完全な例
次の例では、テーブルの作成、データの書き込み、行の追加、読み取り、フィルターとプロジェクションプッシュダウンの適用といった、完全なワークフローを示します。
from __future__ import annotations
import os
from datetime import datetime
import daft
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
from pypaimon.daft import read_paimon, write_paimon
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:
catalog = CatalogFactory.create(CATALOG_OPTIONS)
catalog.create_database(DATABASE, True)
# 1. テーブルの作成 (追加テーブル + unaware bucket)
table_name = f"{DATABASE}.test_daft_paimon_" + datetime.now().strftime("%Y%m%d_%H%M%S")
pa_schema = pa.schema([
pa.field("id", pa.int32(), nullable=False),
pa.field("name", pa.string()),
pa.field("score", pa.int64()),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
options={"bucket": "-1", "file.format": "parquet"},
)
catalog.create_table(table_name, schema, True)
print("created:", table_name)
# 2. 初期バッチの書き込み
initial_df = daft.from_pydict({
"id": [1, 2, 3],
"name": ["Alice", "Bob", "Carol"],
"score": [90, 85, 88],
})
write_paimon(initial_df, table_name, CATALOG_OPTIONS, mode="append")
# 3. 行の追加
append_df = daft.from_pydict({
"id": [4, 5],
"name": ["Dan", "Eve"],
"score": [77, 92],
})
write_paimon(append_df, table_name, CATALOG_OPTIONS, mode="append")
# 4. テーブルの読み取り
df = read_paimon(table_name, CATALOG_OPTIONS).sort("id")
df.show()
assert df.to_pydict()["id"] == [1, 2, 3, 4, 5]
# 5. 述語プッシュダウンとプロジェクションプッシュダウンの適用
filtered = (
read_paimon(table_name, CATALOG_OPTIONS)
.where(daft.col("score") >= 88)
.select("id", "name")
.sort("id")
)
filtered.show()
assert filtered.to_pydict()["id"] == [1, 3, 5]
print("daft + dlf + paimon: ok")
if __name__ == "__main__":
main()
期待される出力:
created: default.test_daft_paimon_20260524_xxxxxx
╭───────┬────────┬───────╮
│ id ┆ name ┆ score │
│ Int32 ┆ String ┆ Int64 │
╞═══════╪════════╪═══════╡
│ 1 ┆ Alice ┆ 90 │
│ 2 ┆ Bob ┆ 85 │
│ 3 ┆ Carol ┆ 88 │
│ 4 ┆ Dan ┆ 77 │
│ 5 ┆ Eve ┆ 92 │
╰───────┴────────┴───────╯
╭───────┬────────╮
│ id ┆ name │
│ Int32 ┆ String │
╞═══════╪════════╡
│ 1 ┆ Alice │
│ 3 ┆ Carol │
│ 5 ┆ Eve │
╰───────┴────────╯
daft + dlf + paimon: ok