本文介绍如何使用 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)