全部產品
Search
文件中心

Data Lake Formation:PyIceberg訪問 DLF

更新時間:Jul 10, 2026

本文介紹如何在EMR on ECS叢集環境中安裝 PyIceberg,並配置其通過 Iceberg REST 協議訪問DLF。

環境準備與安裝

安裝 pyiceberg-dlf

pyiceberg-dlf 是適配 DLF 的 PyIceberg 發行版,發行到 PyPI,包含 REST sigv4 修複、vended OSS 儲存憑證以及憑證到期前自動重新整理能力。

環境要求:Python 3.10 及以上版本。

python3 -m venv venv
source venv/bin/activate
pip install -U pip
# 卸載官方版(與 DLF 適配版不能共存)
pip uninstall -y pyiceberg
# rest-sigv4 為必選項(安裝 boto3,用於 REST sigv4 簽名)
pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"
說明
  • 包名與 import 名不同:包名為 pyiceberg-dlf,import 名仍為 pyiceberg

  • DLF 適配版與官方 pyiceberg 不能共存pyiceberg-dlf 與官方 pyiceberg 都提供 import pyiceberg,後裝的會覆蓋前者。建議使用獨立虛擬環境,或安裝前執行 pip uninstall -y pyiceberg

  • rest-sigv4 不可省略:該 extra 會安裝 boto3 等簽名依賴,缺少時 load_catalog() 將報 ModuleNotFoundError: No module named 'boto3'

  • pandas 按需添加:樣本中使用了 scan.to_pandas(),因此需帶上 pandas extra。

  • pyarrow 須低於 22pyiceberg-dlf[pyarrow] 已將 pyarrow 鎖定到 21.x。pyarrow 22 及以上版本上傳 OSS 時會報 aws-chunked encoding is not supported

程式碼範例

以下 Python 指令碼示範了完整的生命週期操作:串連 Catalog、建立表、寫入資料、讀取資料以及刪除表。

請將指令碼中的預留位置替換為實際配置。

  • ${regionId}:DLF網域名稱,如cn-hangzhou。詳情請參見Iceberg REST服務存取點

  • ${catalogName}:DLF Catalog 名稱。

  • ${accessKey}:AccessKey ID。

  • ${accessKeySecret}:AccessKey Secret。

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. 配置與串連 Catalog
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": "${accessKeyId}",
        "client.secret-access-key": "${accessKeySecret}",
        "client.region": "${regionId}",
    },
)
# ---------------------------------------------------
# 2. 定義表結構與中繼資料
# ---------------------------------------------------
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. 執行測試流程
# ---------------------------------------------------
# 清理舊錶(如果存在)
try:
    catalog.drop_table(identifier=TEST_TABLE_IDENTIFIER)
    print("Existing table dropped.")
except NoSuchTableError:
    print("No existing table to drop.")
# 建立新表
try:
    catalog.create_table(
        identifier=TEST_TABLE_IDENTIFIER,
        schema=TEST_TABLE_SCHEMA,
        partition_spec=TEST_TABLE_PARTITION_SPEC,
        properties=TEST_TABLE_PROPERTIES,
    )
    print("Table created.")
except TableAlreadyExistsError:
    print("Table already exists, will append data.")
# 載入表
table = catalog.load_table(identifier=TEST_TABLE_IDENTIFIER)
print(f"Loaded table: {table}")
# 構造 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,
)
# 寫入資料
print(f"Inserting {data.num_rows} rows...")
table.append(data)
print("Insert finished.")
# 讀取並展示資料
scan = table.scan()
df = scan.to_pandas()
print("First 10 rows via to_pandas():")
print(df.head(10))
# 清理測試表
try:
    catalog.drop_table(identifier=TEST_TABLE_IDENTIFIER)
    print("Table dropped after test.")
except NoSuchTableError:
    print("Table already dropped.")

運行結果展示:

(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#