All Products
Search
Document Center

Data Lake Formation:Akses DLF dengan PyIceberg

Last Updated:Jul 10, 2026

Topik ini menjelaskan cara menginstal PyIceberg di kluster EMR on ECS dan mengonfigurasinya untuk mengakses DLF menggunakan protokol Iceberg REST.

Prasyarat dan instalasi

Instal pyiceberg-dlf

pyiceberg-dlf adalah distribusi PyIceberg yang kompatibel dengan DLF dan dipublikasikan ke PyPI. Distribusi ini mencakup perbaikan REST sigv4, kredensial penyimpanan OSS yang disediakan, serta pembaruan kredensial otomatis sebelum kedaluwarsa.

Persyaratan: Python 3.10 atau versi lebih baru.

python3 -m venv venv
source venv/bin/activate
pip install -U pip
# Uninstal paket resmi (tidak dapat berjalan bersamaan dengan pyiceberg-dlf)
pip uninstall -y pyiceberg
# rest-sigv4 diperlukan (menginstal boto3 untuk penandatanganan REST sigv4)
pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"
Catatan
  • Nama paket berbeda dari nama impor: Nama paketnya adalah pyiceberg-dlf, tetapi nama impornya tetap pyiceberg.

  • pyiceberg-dlf dan pyiceberg resmi tidak dapat berjalan bersamaan: Kedua paket menyediakan import pyiceberg. Menginstal salah satu akan diam-diam menimpa yang lain. Gunakan lingkungan virtual khusus atau jalankan pip uninstall -y pyiceberg sebelum menginstal.

  • rest-sigv4 diperlukan: Komponen tambahan ini menginstal boto3 dan dependensi penandatanganan lainnya. Tanpa komponen ini, load_catalog() akan memunculkan error ModuleNotFoundError: No module named 'boto3'.

  • pandas bersifat opsional: Contoh ini menggunakan scan.to_pandas(), sehingga sertakan komponen tambahan pandas jika Anda memerlukan output dalam bentuk DataFrame.

  • pyarrow harus di bawah versi 22: pyiceberg-dlf[pyarrow] membatasi pyarrow pada versi 21.x. Versi pyarrow 22 ke atas menyebabkan error aws-chunked encoding is not supported saat mengunggah ke OSS.

Contoh kode

Skrip Python berikut menunjukkan cara menghubungkan ke katalog, membuat tabel, menulis data, membaca data, lalu menghapus tabel tersebut.

Ganti placeholder dalam skrip dengan kredensial dan konfigurasi spesifik Anda.

  • ${regionId}: Wilayah layanan DLF Anda, misalnya cn-hangzhou. Untuk informasi selengkapnya, lihat Titik akhir layanan Iceberg REST.

  • ${catalogName}: Nama katalog DLF Anda.

  • ${accessKey}: ID AccessKey Anda.

  • ${accessKeySecret}: Rahasia AccessKey Anda.

import pyarrow as pa
from pyiceberg.catalog import load_catalog
from pyiceberg.exceptions import TableAlreadyExistsError, NoSuchTableError
from pyiceberg.io.pyarrow import schema_to_pyarrow
from pyiceberg.partitioning import PartitionField, PartitionSpec
from pyiceberg.schema import Schema
from pyiceberg.transforms import IdentityTransform
from pyiceberg.types import LongType, NestedField
# 1. Konfigurasikan dan hubungkan ke katalog
catalog = load_catalog(
    "default",
    **{
        "type": "rest",
        "uri": "http://${regionId}-vpc.dlf.aliyuncs.com/iceberg",
        "warehouse": "${catalogName}",
        "rest.signing-name": "DlfNext",
        "rest.signing-region": "${regionId}",
        "rest.sigv4-enabled": "true",
        "client.access-key-id": "${accessKey}",
        "client.secret-access-key": "${accessKeySecret}",
        "client.region": "${regionId}",
    },
)
# ---------------------------------------------------
# 2. Definisikan skema tabel dan metadata
# ---------------------------------------------------
TEST_TABLE_SCHEMA = Schema(
    NestedField(1, "x", LongType(), required=True),
    NestedField(2, "y", LongType(), doc="comment", required=True),
    NestedField(3, "z", LongType(), required=True),
)
TEST_TABLE_IDENTIFIER = ("default", "my_table")
TEST_TABLE_PARTITION_SPEC = PartitionSpec(
    PartitionField(name="x", transform=IdentityTransform(), source_id=1, field_id=1000)
)
TEST_TABLE_PROPERTIES = {"read.split.target.size": "134217728"}  # 128MB
# ---------------------------------------------------
# 3. Jalankan prosedur pengujian
# ---------------------------------------------------
# Hapus tabel lama jika ada
try:
    catalog.drop_table(identifier=TEST_TABLE_IDENTIFIER)
    print("Tabel yang ada telah dihapus.")
except NoSuchTableError:
    print("Tidak ada tabel yang perlu dihapus.")
# Buat tabel baru
try:
    catalog.create_table(
        identifier=TEST_TABLE_IDENTIFIER,
        schema=TEST_TABLE_SCHEMA,
        partition_spec=TEST_TABLE_PARTITION_SPEC,
        properties=TEST_TABLE_PROPERTIES,
    )
    print("Tabel berhasil dibuat.")
except TableAlreadyExistsError:
    print("Tabel sudah ada, data akan ditambahkan.")
# Muat tabel
table = catalog.load_table(identifier=TEST_TABLE_IDENTIFIER)
print(f"Tabel dimuat: {table}")
# Buat tabel PyArrow
arrow_schema = schema_to_pyarrow(table.schema())
data = pa.Table.from_pydict(
    {
        "x": [1, 2, 3],
        "y": [10, 20, 30],
        "z": [100, 200, 300],
    },
    schema=arrow_schema,
)
# Tulis data
print(f"Menyisipkan {data.num_rows} baris...")
table.append(data)
print("Penyisipan selesai.")
# Baca dan tampilkan data
scan = table.scan()
df = scan.to_pandas()
print("10 baris pertama melalui to_pandas():")
print(df.head(10))
# Bersihkan tabel uji
try:
    catalog.drop_table(identifier=TEST_TABLE_IDENTIFIER)
    print("Tabel dihapus setelah pengujian.")
except NoSuchTableError:
    print("Tabel sudah dihapus.")

Contoh output:

(myenv) root@iZbp1h4zr65vjcvz9w3080Z:~/workspace# python test.py
Existing table dropped.
Table created.
Loaded table: my_table(
    1: x: required long,
    2: y: required long (comment),
    3: z: required long
),
partition by: [x],
sort order: [],
snapshot: null
Inserting 3 rows...
Insert finished.
First 10 rows via to_pandas():
   x   y    z
0  1  10  100
1  2  20  200
2  3  30  300
Table dropped after test.
(myenv) root@iZbp1h4zr65vjcvz9w3080Z:~/workspace# pip list|grep -E "pyarrow|pyiceberg|boto3|pandas"
boto3                 1.42.15
pandas                2.3.3
pyarrow               19.0.0
pyiceberg             0.10.0.dev0
(myenv) root@iZbp1h4zr65vjcvz9w3080Z:~/workspace#