Daft は高性能な分散データフレームエンジンです。Daft を使用して、Alibaba Cloud Data Lake Formation (DLF) 上の Apache Iceberg テーブルに対してデータの読み書きを行い、フィルタリングや集計などのデータフレーム操作を実行できます。
前提条件
DLF Iceberg REST サービスは VPC 内からのみアクセス可能です。このトピックのコードは、DLF と同じリージョンの VPC 環境 (ECS インスタンスや EMR クラスターなど) 内から実行してください。各リージョンのエンドポイントについては、「Iceberg REST エンドポイント」をご参照ください。
依存関係のインストール
Python 3.10 以降。
DLF 互換の PyIceberg パッケージ (
pyiceberg-dlf) をインストールします。python3 -m venv venv source venv/bin/activate pip install -U pip # pyiceberg をアンインストール (pyiceberg-dlf と共存できないため) pip uninstall -y pyiceberg # rest-sigv4 が必要 (REST sigv4 署名用に boto3 をインストール) pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"Daft をインストールします。
pip install "daft>=0.7.17"
パラメータの設定
以下の情報を準備してください。
パラメータ | 説明 |
| AccessKey ID。 |
| AccessKey secret。 |
| DLF がデプロイされているリージョン ID です。例: |
| DLF のカタログ名 (Iceberg ウェアハウスに対応)。 |
| 対象データベースの名前 (Iceberg 名前空間)。 |
AccessKey 認証情報を保護してください。コードにハードコードしたり、コードリポジトリにコミットしたりしないでください。環境変数またはシークレット管理サービスから取得することを推奨します。
接続と初期化
DLF カタログへの接続
PyIceberg の load_catalog を使用して、Iceberg REST プロトコル経由で DLF カタログに接続します。
from pyiceberg.catalog import load_catalog
REGION = "${regionId}"
catalog = load_catalog(
"dlf",
**{
"type": "rest",
"uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
"warehouse": "${catalogName}",
"rest.signing-name": "DlfNext",
"rest.signing-region": REGION,
"rest.sigv4-enabled": "true",
"client.access-key-id": "${accessKeyId}",
"client.secret-access-key": "${accessKeySecret}",
"client.region": REGION,
"s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com",
},
)テーブルの作成または読み込み
新しいテーブルの作成
PyIceberg を使用して Iceberg テーブルを作成します。DLF がテーブルメタデータを管理します。
from pyiceberg.schema import Schema
from pyiceberg.types import LongType, StringType, DoubleType, NestedField
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC
schema = Schema(
NestedField(field_id=1, name="id", field_type=LongType(), required=True),
NestedField(field_id=2, name="name", field_type=StringType(), required=False),
NestedField(field_id=3, name="category", field_type=StringType(), required=False),
NestedField(field_id=4, name="value", field_type=DoubleType(), required=False),
)
table = catalog.create_table(
identifier=("${database}", "daft_demo_table"),
schema=schema,
partition_spec=UNPARTITIONED_PARTITION_SPEC,
)
print("Table created. Data location:", table.location())既存のテーブルの読み込み
既存のテーブルを操作するには、catalog.load_table() を呼び出します。
table = catalog.load_table(("${database}", "table_name"))データ操作
データの書き込み
Daft を使用して DataFrame を構築し、Iceberg テーブルに書き込みます。 write_iceberg は、この操作で書き込まれたデータファイルの概要テーブルを返します。
Daft 0.7.15 では、oss:// パス上の Iceberg テーブルに対する OSS アクセスが自動的に設定されるため、write_iceberg および read_iceberg では手動の IOConfig は不要になりました。
import daft
df = daft.from_pydict({
"id": [1, 2, 3, 4, 5],
"name": ["item_1", "item_2", "item_3", "item_4", "item_5"],
"category": ["A", "B", "A", "B", "A"],
"value": [10.5, 21.0, 31.5, 42.0, 52.5],
})
result = df.write_iceberg(table, mode="append")
result.show()mode パラメーターは "append" (データの追記) と "overwrite" (既存のデータの上書き) に対応しています。
データの読み込み
データの書き込み後、テーブルを再読み込みして最新のスナップショットを取得し、STS 認証情報を更新してから、Daft を使用してデータを読み込みます。
table = catalog.load_table(("${database}", "daft_demo_table"))
df = daft.read_iceberg(table)
df.show()daft.read_iceberg は遅延であり、DataFrame ハンドルを返して実行計画を構築するため、データは show() や collect() などのアクションを呼び出したときにのみ OSS から読み取られます。
データの変換
テーブルから読み込んだデータフレーム (5 行のサンプルデータを含む) に対して、フィルタリング、派生列の作成、集計、ソートを実行できます。
df = daft.read_iceberg(table).collect()
# フィルタリング: value が 30 より大きい行
df.where(df["value"] > 30).show()
# 派生列: value_x2 = value * 2 を追加
df.with_column("value_x2", df["value"] * 2).show()
# グループ化と集計: category ごとに行数と value の合計を計算
df.groupby("category").agg(
daft.col("id").count().alias("row_count"),
daft.col("value").sum().alias("value_sum"),
).sort("category").show()
# ソート: value の降順
df.sort("value", desc=True).show()完全な例
以下のコードは、上記のすべてのステップを組み合わせたもので、そのままコピーして実行できます。
import uuid
import daft
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import LongType, StringType, DoubleType, NestedField
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC
# ==================== 設定 ====================
ACCESS_KEY_ID = "${accessKeyId}"
ACCESS_KEY_SECRET = "${accessKeySecret}"
REGION = "${regionId}"
CATALOG_NAME = "${catalogName}"
DATABASE = "${database}"
# =======================================================
def create_catalog():
"""DLF Iceberg REST カタログに接続します。"""
return load_catalog(
"dlf",
**{
"type": "rest",
"uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
"warehouse": CATALOG_NAME,
"rest.signing-name": "DlfNext",
"rest.signing-region": REGION,
"rest.sigv4-enabled": "true",
"client.access-key-id": ACCESS_KEY_ID,
"client.secret-access-key": ACCESS_KEY_SECRET,
"client.region": REGION,
"s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com",
},
)
def create_table(catalog, table_name):
"""PyIceberg を使用してパーティション化されていない Iceberg テーブルを作成します。"""
schema = Schema(
NestedField(field_id=1, name="id", field_type=LongType(), required=True),
NestedField(field_id=2, name="name", field_type=StringType(), required=False),
NestedField(field_id=3, name="category", field_type=StringType(), required=False),
NestedField(field_id=4, name="value", field_type=DoubleType(), required=False),
)
return catalog.create_table(
identifier=(DATABASE, table_name),
schema=schema,
partition_spec=UNPARTITIONED_PARTITION_SPEC,
)
def main():
catalog = create_catalog()
table_name = f"daft_demo_{uuid.uuid4().hex[:8]}"
table = None
try:
# 1. テーブルの作成 (PyIceberg)
table = create_table(catalog, table_name)
print("Table created. Data location:", table.location())
# 2. Daft を使用したデータの書き込み
df = daft.from_pydict({
"id": [1, 2, 3, 4, 5],
"name": ["item_1", "item_2", "item_3", "item_4", "item_5"],
"category": ["A", "B", "A", "B", "A"],
"value": [10.5, 21.0, 31.5, 42.0, 52.5],
})
df.write_iceberg(table, mode="append").show()
# 3. テーブルを再読み込みして最新のスナップショットを取得し、Daft を使用して読み込み
table = catalog.load_table((DATABASE, table_name))
result = daft.read_iceberg(table).collect()
result.show()
# 4. データフレーム変換
result.groupby("category").agg(
daft.col("id").count().alias("row_count"),
daft.col("value").sum().alias("value_sum"),
).sort("category").show()
finally:
# 5. テーブルの削除 (PyIceberg)
if table is not None:
catalog.drop_table((DATABASE, table_name))
print("Table dropped:", table_name)
if __name__ == "__main__":
main()書き込みステップの期待される出力(write_iceberg 結果テーブル)は、次のようになります。
╭───────────┬───────┬───────────┬────────────────────────────────╮
│ operation ┆ rows ┆ file_size ┆ file_name │
╞═══════════╪═══════╪═══════════╪════════════════════════════════╡
│ ADD ┆ 5 ┆ 1708 ┆ oss://<bucket>/.../xxx.parquet │
╰───────────┴───────┴───────────┴────────────────────────────────╯付録
用語
コンポーネント | 説明 |
DLF | Data Lake Formation。Iceberg REST カタログサービスと統合されたテーブルメタデータ管理を提供します。 |
PyIceberg | Apache Iceberg の Python クライアント。DLF カタログに接続し、テーブルの作成と削除、トランザクションのコミットを実行します。 |
Daft | 分散データフレームエンジン。Iceberg テーブルへのデータの書き込みと読み込み、およびデータフレーム変換を実行します。 |
OSS | Object Storage Service。Iceberg テーブルデータファイル (Parquet 形式) の物理ストレージレイヤー。 |
アーキテクチャ
Daft と PyIceberg は、明確な役割分担のもとで連携します。
コントロールプレーン (PyIceberg):Iceberg REST プロトコルを使用して DLF に接続します。カタログ接続、テーブルの作成と削除、スナップショットのコミットを処理します。
データプレーン (Daft):Iceberg テーブルのデータファイルは OSS に保存されます。Daft はこれらの Parquet ファイルを並列に読み書きし、データフレーム計算機能 (フィルタリング、派生列、集計、ソートなど) を提供します。
書き込み時、Daft は Parquet データファイルを書き込み、PyIceberg を介して Iceberg スナップショットをアトミックにコミットします。読み込み時、Daft は Iceberg メタデータを活用してパーティションプルーニングとファイルフィルタリングを実行します。