Daft adalah mesin DataFrame terdistribusi berkinerja tinggi. Anda dapat menggunakan Daft untuk membaca dan menulis tabel Apache Iceberg di Alibaba Cloud Data Lake Formation (DLF), serta melakukan operasi DataFrame seperti filtering dan agregasi.
Prasyarat
Layanan REST Iceberg DLF hanya dapat diakses dari dalam VPC. Jalankan kode dalam topik ini dari lingkungan VPC di Wilayah yang sama dengan DLF (misalnya Instance ECS atau kluster EMR). Untuk titik akhir setiap wilayah, lihat Titik akhir REST Iceberg.
Instal dependensi
Python 3.10 atau versi lebih baru.
Instal paket PyIceberg yang kompatibel dengan DLF (
pyiceberg-dlf).python3 -m venv venv source venv/bin/activate pip install -U pip # Uninstal pyiceberg, yang tidak dapat berdampingan dengan pyiceberg-dlf pip uninstall -y pyiceberg # rest-sigv4 diperlukan (menginstal boto3 untuk penandatanganan REST sigv4) pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"Instal Daft.
pip install "daft>=0.7.17"
Konfigurasi parameter
Siapkan informasi berikut:
Parameter | Deskripsi |
| ID AccessKey Anda. |
| Rahasia AccessKey Anda. |
| ID Wilayah tempat DLF dideploy, misalnya |
| Nama katalog di DLF (berkorespondensi dengan gudang Iceberg). |
| Nama database target (namespace Iceberg). |
Lindungi kredensial AccessKey Anda. Jangan hardcode di kode atau commit ke repositori kode. Kami merekomendasikan mengambilnya dari Variabel lingkungan atau layanan manajemen rahasia.
Koneksi dan inisialisasi
Koneksi ke katalog DLF
Koneksikan ke katalog DLF melalui protokol REST Iceberg dengan menggunakan load_catalog dari PyIceberg.
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",
},
)Buat atau muat tabel
Buat tabel baru
Buat tabel Iceberg menggunakan PyIceberg. DLF mengelola metadata tabel.
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("Tabel dibuat. Lokasi data:", table.location())Muat tabel yang sudah ada
Untuk bekerja dengan tabel yang sudah ada, panggil catalog.load_table().
table = catalog.load_table(("${database}", "table_name"))Operasi data
Tulis data
Buat DataFrame menggunakan Daft dan tulis ke tabel Iceberg. write_iceberg mengembalikan tabel ringkasan file data yang ditulis dalam operasi ini.
Daft 0.7.15 secara otomatis mengonfigurasi akses OSS untuk tabel Iceberg pada path oss://, sehingga write_iceberg dan read_iceberg tidak memerlukan IOConfig manual.
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()Parameter mode mendukung "append" (tambahkan data) dan "overwrite" (timpa data yang ada).
Baca data
Setelah menulis data, muat ulang tabel untuk mendapatkan Snapshot terbaru dan refresh kredensial STS, lalu baca menggunakan Daft.
table = catalog.load_table(("${database}", "daft_demo_table"))
df = daft.read_iceberg(table)
df.show()daft.read_iceberg bersifat lazy: mengembalikan handle DataFrame dan menyusun rencana eksekusi. Data hanya dibaca dari OSS saat Anda memanggil aksi seperti show() atau collect().
Transformasi data
Dengan DataFrame yang dibaca dari tabel (berisi lima baris contoh), Anda dapat melakukan filtering, pembuatan kolom turunan, agregasi, dan Pengurutan.
df = daft.read_iceberg(table).collect()
# Filter: baris dengan nilai lebih besar dari 30
df.where(df["value"] > 30).show()
# Kolom turunan: tambahkan value_x2 = value * 2
df.with_column("value_x2", df["value"] * 2).show()
# Kelompokkan dan agregasi: hitung jumlah baris dan jumlah nilai berdasarkan kategori
df.groupby("category").agg(
daft.col("id").count().alias("row_count"),
daft.col("value").sum().alias("value_sum"),
).sort("category").show()
# Urutkan: menurun berdasarkan nilai
df.sort("value", desc=True).show()Contoh lengkap
Kode berikut menggabungkan semua langkah sebelumnya dan dapat disalin serta dijalankan langsung.
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
# ==================== Konfigurasi ====================
ACCESS_KEY_ID = "${accessKeyId}"
ACCESS_KEY_SECRET = "${accessKeySecret}"
REGION = "${regionId}"
CATALOG_NAME = "${catalogName}"
DATABASE = "${database}"
# =======================================================
def create_catalog():
"""Koneksikan ke Katalog REST Iceberg DLF."""
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):
"""Buat tabel Iceberg tanpa partisi menggunakan PyIceberg."""
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. Buat tabel (PyIceberg)
table = create_table(catalog, table_name)
print("Tabel dibuat. Lokasi data:", table.location())
# 2. Tulis data menggunakan 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. Muat ulang tabel untuk mendapatkan Snapshot terbaru, lalu baca menggunakan Daft
table = catalog.load_table((DATABASE, table_name))
result = daft.read_iceberg(table).collect()
result.show()
# 4. Transformasi DataFrame
result.groupby("category").agg(
daft.col("id").count().alias("row_count"),
daft.col("value").sum().alias("value_sum"),
).sort("category").show()
finally:
# 5. Hapus tabel (PyIceberg)
if table is not None:
catalog.drop_table((DATABASE, table_name))
print("Tabel dihapus:", table_name)
if __name__ == "__main__":
main()Output yang diharapkan dari langkah penulisan (tabel hasil write_iceberg) tampak seperti berikut:
╭───────────┬───────┬───────────┬────────────────────────────────╮
│ operation ┆ rows ┆ file_size ┆ file_name │
╞═══════════╪═══════╪═══════════╪════════════════════════════════╡
│ ADD ┆ 5 ┆ 1708 ┆ oss://<bucket>/.../xxx.parquet │
╰───────────┴───────┴───────────┴────────────────────────────────╯Lampiran
Terminologi
Komponen | Deskripsi |
DLF | Data Lake Formation. Menyediakan layanan Katalog REST Iceberg dan manajemen metadata tabel terpadu. |
PyIceberg | Klien Python untuk Apache Iceberg. Terhubung ke katalog DLF, membuat dan menghapus tabel, serta melakukan commit transaksi. |
Daft | Mesin DataFrame terdistribusi. Menulis dan membaca data dari tabel Iceberg, serta melakukan transformasi DataFrame. |
OSS | Object Storage Service. Lapisan penyimpanan fisik untuk file data tabel Iceberg (format Parquet). |
Arsitektur
Daft dan PyIceberg bekerja sama dengan pemisahan tanggung jawab yang jelas:
Control plane (PyIceberg): Terhubung ke DLF menggunakan protokol REST Iceberg. Menangani koneksi katalog, pembuatan dan penghapusan tabel, serta commit snapshot.
Data plane (Daft): File data tabel Iceberg disimpan di OSS. Daft membaca dan menulis file Parquet ini secara paralel serta menyediakan kemampuan komputasi DataFrame (filtering, kolom turunan, agregasi, Pengurutan, dll.).
Untuk penulisan, Daft menulis file data Parquet dan melakukan commit snapshot Iceberg secara atomik melalui PyIceberg. Untuk pembacaan, Daft memanfaatkan metadata Iceberg untuk pemangkasan partisi dan filtering file.