Daft 通過 pypaimon 的 read_paimon / write_paimon 介面讀寫 DLF Catalog 管理的 Paimon 表,支援謂詞與列裁剪下推、時間旅行查詢。本文給出最小可啟動並執行端到端樣本。
如需操作 BLOB / list[binary] 等多模態列或執行多模態 ETL,請參見基於Daft和DLF的自動駕駛多模態資料湖實踐。
前提條件
步驟一:安裝依賴
下載pypaimon離線安裝包pypaimon-1.5.dev20260608.tar.gz
安裝 Python 依賴
在pypaimon離線安裝包所在目錄執行以下命令安裝所有依賴:
pip install daft pyarrow pandas requests pypaimon-1.5.dev20260608.tar.gz驗證環境
執行以下代碼驗證環境是否就緒:
import daft from pypaimon.daft import read_paimon, write_paimon print("daft:", daft.__version__, "pypaimon.daft: OK")
步驟二:配置 Catalog 串連
通過以下配置串連 DLF Catalog:
CATALOG_OPTIONS = {
"metastore": "rest",
"uri": "http://<DLF-ENDPOINT>", # 公網訪問須使用 HTTPS
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": "<ACCESS-KEY-ID>",
"dlf.access-key-secret": "<ACCESS-KEY-SECRET>",
"dlf.oss-endpoint": "<OSS-ENDPOINT>", # VPC 可選,公網必填
}參數說明如下:
參數 | 說明 |
metastore | 固定為 rest。 |
uri | DLF Paimon REST 訪問端點,詳見服務存取點。VPC 訪問支援 HTTP 和 HTTPS 協議,公網訪問必須使用 HTTPS 協議。VPC訪問樣本:http://ap-southeast-1-vpc.dlf.aliyuncs.com |
warehouse | DLF Catalog 名稱。 |
token.provider | AccessKey 認證模式下填 dlf。 |
dlf.region | DLF 所屬地區 ID(如 cn-hangzhou),詳見服務存取點。 |
dlf.access-key-id / dlf.access-key-secret | DLF AccessKey。 |
dlf.security-token | (可選)STS 安全性權杖。 |
dlf.oss-endpoint | (可選)OSS 公網訪問 endpoint。 |
寫入資料
建立新表
新表必須先用 PyPaimon 的 create_table() 建立,再用 Daft 讀寫:
from datetime import datetime
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
catalog = CatalogFactory.create(CATALOG_OPTIONS)
catalog.create_database("default", True)
table_name = "default.test_daft_paimon_" + datetime.now().strftime("%Y%m%d_%H%M%S")
pa_schema = pa.schema([
pa.field("id", pa.int32(), nullable=False),
pa.field("name", pa.string()),
pa.field("score", pa.int64()),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
options={"bucket": "-1", "file.format": "parquet"},
)
catalog.create_table(table_name, schema, True)追加寫入已有表
使用 write_paimon 將 Daft DataFrame 寫入已有 Paimon 表:
import daft
from pypaimon.daft import write_paimon
append_df = daft.from_pydict({
"id": [4, 5],
"name": ["Dan", "Eve"],
"score": [77, 92],
})
write_paimon(append_df, "default.your_table", CATALOG_OPTIONS, mode="append")mode 支援以下兩種取值:
"append":追加寫入。"overwrite":覆蓋寫入。
讀取已有表
基本讀取
使用 read_paimon 讀取指定表的全部資料:
from pypaimon.daft import read_paimon
df = read_paimon("default.your_table", CATALOG_OPTIONS)
df.show()謂詞與列裁剪
謂詞與列裁剪會自動下推到 Paimon 掃描層:
import daft
df = (
read_paimon("default.your_table", CATALOG_OPTIONS)
.where(daft.col("score") >= 88)
.select("id", "name")
)
df.show()時間旅行
支援通過 snapshot_id 或 tag_name 讀取歷史版本:
# 通過 snapshot_id
df = read_paimon("default.your_table", CATALOG_OPTIONS, snapshot_id=42)
# 通過 tag_name
df = read_paimon("default.your_table", CATALOG_OPTIONS, tag_name="v1")完整樣本
以下樣本示範了完整的端到端流程:建立新表 → Daft 寫入 → 追加 → 讀取 → 過濾投影。
from __future__ import annotations
import os
from datetime import datetime
import daft
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
from pypaimon.daft import read_paimon, write_paimon
CATALOG_OPTIONS = {
"metastore": "rest",
"uri": "http://<DLF-ENDPOINT>",
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": os.environ["DLF_ACCESS_KEY_ID"],
"dlf.access-key-secret": os.environ["DLF_ACCESS_KEY_SECRET"],
"dlf.oss-endpoint": "<OSS-ENDPOINT>",
}
DATABASE = "default"
def main() -> None:
catalog = CatalogFactory.create(CATALOG_OPTIONS)
catalog.create_database(DATABASE, True)
# 1. 建表(append 表 + unaware bucket)
table_name = f"{DATABASE}.test_daft_paimon_" + datetime.now().strftime("%Y%m%d_%H%M%S")
pa_schema = pa.schema([
pa.field("id", pa.int32(), nullable=False),
pa.field("name", pa.string()),
pa.field("score", pa.int64()),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
options={"bucket": "-1", "file.format": "parquet"},
)
catalog.create_table(table_name, schema, True)
print("created:", table_name)
# 2. Daft 寫入首批資料
initial_df = daft.from_pydict({
"id": [1, 2, 3],
"name": ["Alice", "Bob", "Carol"],
"score": [90, 85, 88],
})
write_paimon(initial_df, table_name, CATALOG_OPTIONS, mode="append")
# 3. Daft 追加
append_df = daft.from_pydict({
"id": [4, 5],
"name": ["Dan", "Eve"],
"score": [77, 92],
})
write_paimon(append_df, table_name, CATALOG_OPTIONS, mode="append")
# 4. Daft 讀取
df = read_paimon(table_name, CATALOG_OPTIONS).sort("id")
df.show()
assert df.to_pydict()["id"] == [1, 2, 3, 4, 5]
# 5. 謂詞 + 投影下推
filtered = (
read_paimon(table_name, CATALOG_OPTIONS)
.where(daft.col("score") >= 88)
.select("id", "name")
.sort("id")
)
filtered.show()
assert filtered.to_pydict()["id"] == [1, 3, 5]
print("daft + dlf + paimon: ok")
if __name__ == "__main__":
main()預期輸出如下:
created: default.test_daft_paimon_20260524_xxxxxx
╭───────┬────────┬───────╮
│ id ┆ name ┆ score │
│ Int32 ┆ String ┆ Int64 │
╞═══════╪════════╪═══════╡
│ 1 ┆ Alice ┆ 90 │
│ 2 ┆ Bob ┆ 85 │
│ 3 ┆ Carol ┆ 88 │
│ 4 ┆ Dan ┆ 77 │
│ 5 ┆ Eve ┆ 92 │
╰───────┴────────┴───────╯
╭───────┬────────╮
│ id ┆ name │
│ Int32 ┆ String │
╞═══════╪════════╡
│ 1 ┆ Alice │
│ 3 ┆ Carol │
│ 5 ┆ Eve │
╰───────┴────────╯
daft + dlf + paimon: ok