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]"Nama paket berbeda dari nama impor: Nama paketnya adalah
pyiceberg-dlf, tetapi nama impornya tetappyiceberg.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 jalankanpip uninstall -y pyicebergsebelum menginstal.rest-sigv4 diperlukan: Komponen tambahan ini menginstal boto3 dan dependensi penandatanganan lainnya. Tanpa komponen ini,
load_catalog()akan memunculkan errorModuleNotFoundError: No module named 'boto3'.pandas bersifat opsional: Contoh ini menggunakan
scan.to_pandas(), sehingga sertakan komponen tambahanpandasjika 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 erroraws-chunked encoding is not supportedsaat 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, misalnyacn-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#