All Products
Search
Document Center

Data Lake Formation:Gunakan Daft dengan tabel Iceberg DLF

Last Updated:Jul 10, 2026

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

Catatan

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

  1. Python 3.10 atau versi lebih baru.

  2. 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]"
  3. Instal Daft.

    pip install "daft>=0.7.17"

Konfigurasi parameter

Siapkan informasi berikut:

Parameter

Deskripsi

${accessKeyId}

ID AccessKey Anda.

${accessKeySecret}

Rahasia AccessKey Anda.

${regionId}

ID Wilayah tempat DLF dideploy, misalnya ap-southeast-1. Untuk daftar ID wilayah, lihat Titik akhir dan akses jaringan publik.

${catalogName}

Nama katalog di DLF (berkorespondensi dengan gudang Iceberg).

${database}

Nama database target (namespace Iceberg).

Penting

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.

Catatan

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.