本文介紹如何使用 PyPaimon 讀取 DLF 多模態表中的資料,涵蓋單機讀取和 Ray 分布式讀取等多種情境。
前提條件
使用 PyPaimon 讀取多模態表前,請確保已安裝以下依賴:
pip install pyjindosdk
pip install pypaimon-1.5.dev20260727.tar.gz
-
擷取PyPaimon安裝包:pypaimon-1.5.dev20260727.tar.gz
-
建議生產環境安裝 pyjindosdk,通過 pyjindosdk 讀取 OSS 可以避免回退到 legacy OSS 路徑,從而在高並發讀取時保持穩定的效能。
串連初始化
使用 PyPaimon 讀取多模態表前,需要先串連 Catalog 並擷取表對象。後續所有讀取情境(單機讀取、Ray 分布式讀取、Daft on Ray)均複用此 table 對象。
import pypaimon.multimodal as pm
catalog_options = {
"metastore": "rest",
"uri": "<DLF_ENDPOINT>",
"warehouse": "<CATALOG_NAME>",
"token.provider": "dlf",
"dlf.region": "<REGION_ID>",
"dlf.oss-endpoint": "<OSS_ENDPOINT>",
"dlf.access-key-id": "<ACCESS_KEY_ID>",
"dlf.access-key-secret": "<ACCESS_KEY_SECRET>",
"dlf.security-token": "<SECURITY_TOKEN>",
}
conn = pm.connect(database="<DATABASE_NAME>", options=catalog_options)
table = conn.get_table("<TABLE_NAME>")
串連參數說明
|
參數 |
說明 |
|
|
Paimon 資料庫名稱。 |
|
|
DLF Endpoint,例如 |
|
|
DLF Catalog 名稱。 |
|
|
地區 ID,例如 |
|
|
OSS Endpoint,例如 |
|
|
阿里雲 AccessKey ID。 |
|
|
阿里雲 AccessKey Secret。 |
|
|
STS 臨時憑證的 Security Token。使用長期 AccessKey 時可省略。 |
|
|
多模態表名稱。 |
單機讀取
少量資料一次性讀取
當您需要讀取單個 clip 或少量 clip,且結果能夠放入記憶體時,可以使用 read_blobs 方法一次性將資料物化到記憶體。
scalar, blobs = (
table.scan()
.where("clip_id = 'xxx'")
.select(["clip_id", "frame_index", "camera_0", "camera_1", "camera_2", "camera_3"])
.read_blobs(
["camera_0", "camera_1", "camera_2", "camera_3"],
parallelism=16,
)
)
大批量流式讀取
當需要讀取大批量多 clip 資料時,推薦使用 stream_blobs 方法分 batch 流式消費,避免將大量 Blob bytes 一次性載入到記憶體。
clip_ids = ["clip_001", "clip_002", "clip_003"]
clip_filter = "clip_id IN (" + ", ".join(f"'{clip_id}'" for clip_id in clip_ids) + ")"
for scalar_batch, blobs in (
table.scan()
.where(clip_filter)
.select(["clip_id", "frame_index", "camera_0", "camera_1", "camera_2", "camera_3"])
.stream_blobs(
["camera_0", "camera_1", "camera_2", "camera_3"],
parallelism=16,
)
):
consume_batch(scalar_batch, blobs)
參數說明
|
參數 |
說明 |
|
|
Blob 讀取線程數。 |
|
|
|
訓練讀取連續 frame range
訓練取數情境下,通常需要從同一個 clip 中讀取連續的 frame window(如 16、32 或 64 幀)。以下樣本展示如何讀取單個 clip 內的一段連續 frame range。
start = 1000
read_frames = 512
for scalar_batch, blobs in (
table.scan()
.where(
f"clip_id = 'xxx' "
f"AND frame_index >= {start} "
f"AND frame_index < {start + read_frames}"
)
.select(["clip_id", "frame_index", "camera_0", "camera_1", "camera_2", "camera_3"])
.stream_blobs(
["camera_0", "camera_1", "camera_2", "camera_3"],
parallelism=16,
)
):
consume_training_range(scalar_batch, blobs)
Ray 分布式讀取 Blob
在 Ray 情境下,推薦分兩步完成分布式讀取:
-
調用
scan().to_ray()分布式讀取標量列和 Blob descriptor。 -
調用
table.map_with_blobs()在 Ray worker 中讀取 Blob bytes,並交給使用者自訂函數(UDF)消費。
基本用法
import ray
import pyarrow as pa
ray.init(address="auto", ignore_reinit_error=True)
clip_ids = ["clip_001", "clip_002", "clip_003"]
scalar_cols = ["clip_id", "frame_index"]
blob_cols = ["camera_0", "camera_1", "camera_2", "camera_3"]
select_cols = scalar_cols + blob_cols
clip_filter = "clip_id IN (" + ", ".join(f"'{clip_id}'" for clip_id in clip_ids) + ")"
def process_batch(scalar_batch, blobs):
"""
scalar_batch: pyarrow.Table,包含 clip_id、frame_index 等標量列。
blobs: dict[str, list[bytes | None]],key 為 blob_cols 中的列名。
傳回值需要是一個小的 pyarrow.Table;只做 side-effect 時可返回空表。
"""
camera_0 = blobs["camera_0"]
# 不建議返回原始 Blob bytes,避免把大量資料物化到 Ray object store。
return pa.table({
"rows": [scalar_batch.num_rows],
})
ds = (
table.scan()
.where(clip_filter)
.select(select_cols)
.to_ray()
)
result_ds = table.map_with_blobs(
ds,
blob_cols,
process_batch,
)
# 觸發執行。Ray Dataset 是 lazy 的,不消費 result_ds 就不會真正讀取 Blob。
for _ in result_ds.iter_batches(batch_format="pyarrow"):
pass
參數說明
|
參數 |
說明 |
|
|
預設由 Ray 根據可用資源和輸入資料規模決定讀取並發和 block 數,通常無需手動設定。 |
|
|
每個 Ray task 內部讀取 Blob bytes 的線程數,不是 Ray worker 數。預設值為 64。 |
|
|
每次傳給 UDF 的行數,預設值為 1024。Blob 較大或 worker 記憶體壓力高時,建議調小到 128、256 或 512。 |
配置 Ray task 重試
大規模冷讀情境下,建議為 map_with_blobs() 配置 Ray task 重試。
result_ds = table.map_with_blobs(
ds,
blob_cols,
process_batch,
parallelism=8,
batch_size=512,
ray_remote_args={
"max_retries": 3,
"retry_exceptions": True,
},
)
如果 descriptor 規模也較大,也可以為 to_ray() 傳入相同的 ray_remote_args。
Daft on Ray 分布式讀取 Blob
Daft on Ray 是另一種分布式讀取路徑,適用於已有 Daft 作業鏈路的情境。
基本用法
import datetime as dt
import daft
from daft import col, runners
import pyarrow as pa
import ray
from pypaimon.daft import read_paimon, read_blob
ray.init(address="auto", ignore_reinit_error=True)
runners.set_runner_ray(address="auto", noop_if_initialized=True)
table_identifier = "<DATABASE_NAME>.<TABLE_NAME>"
clip_ids = ["clip_001", "clip_002", "clip_003"]
scalar_cols = ["clip_id", "frame_index", "collected_date"]
blob_cols = ["camera_0", "camera_1", "camera_2", "camera_3"]
target_date = dt.date(2026, 7, 1)
df = read_paimon(table_identifier, catalog_options)
df = df.where(
(col("collected_date") == target_date)
& col("clip_id").is_in(clip_ids)
)
# read_blob 會把 Daft File/Blob descriptor 列讀取成 binary bytes。
# max_concurrency 是每個 Daft batch UDF 內部讀取 Blob 的並發度。
blob_bytes_cols = []
for blob_col in blob_cols:
bytes_col = blob_col.replace(".", "_") + "_bytes"
blob_bytes_cols.append(bytes_col)
df = df.select(
*[col(name) for name in scalar_cols],
*[
read_blob(
col(blob_col),
catalog_options,
table_identifier,
max_concurrency=16,
).alias(bytes_col)
for blob_col, bytes_col in zip(blob_cols, blob_bytes_cols)
],
)
# 樣本:這裡只統計 bytes 大小。真實業務中可接 decode、preprocess、inference、write 等邏輯。
# 不建議長期保留或寫出原始 Blob bytes,避免大量資料進入 Ray object store。
agg_exprs = [col("clip_id").count().alias("rows")]
for bytes_col in blob_bytes_cols:
agg_exprs.append(col(bytes_col).count().alias(bytes_col + "_count"))
agg_exprs.append(col(bytes_col).str.length().sum().alias(bytes_col + "_sum"))
result = df.agg(*agg_exprs).collect()
使用 open_blob 逐個開啟 Blob 流
如果需要在自訂 Daft UDF 中逐個開啟 Blob 流進行處理,可以使用 open_blob() 方法。其中 max_concurrency 控制每個 Daft batch UDF 內部讀取 Blob 的並發度。
from pypaimon.daft import open_blob
def consume_one_blob(file):
with open_blob(file, catalog_options, table_identifier) as stream:
data = stream.read()
# decode / preprocess / inference / write
return len(data)