演示如何用 Daft + DLF Paimon 表管理自动驾驶多摄像头帧数据:原图与缩略图以 BLOB 形式入湖、按 channel/车速/天气快速过滤、按需拉取大字节、为训练集打 tag 回溯,并接入 PIL/模型预处理 UDF。
场景
智能驾驶车队从多摄像头采集帧序列入湖,下游用于训练集筛选、回归集抽样、感知模型批量预处理。典型需求如下:
按 channel(相机通道)、时间窗、自车状态等维度快速定位帧。
BLOB 形式的图像不能每次都全量物化。
训练集快照(snapshot)可回溯、可打 tag(tag 即 snapshot 的命名引用)。
各组件分工如下:
组件 | 提供能力 |
Paimon | 存储、Schema 演进、时间旅行 |
Daft | lazy DataFrame、优化器、BLOB-as- |
DLF | Catalog 服务、OSS 凭证管理 |
前提条件
安装依赖
下载 pypaimon 离线安装包:pypaimon-1.5.dev20260704.tar.gz
安装 Python 依赖
在离线包所在目录执行以下命令,一次性安装通用依赖与本场景额外依赖:
pip install \ daft \ 'pyarrow>=16,<=21' \ pandas \ requests \ pypaimon-1.5.dev20260608.tar.gz \ modelscope \ datasets \ pillow \ oss2验证环境
import daft from pypaimon.daft import read_paimon, write_paimon print("daft:", daft.__version__, "pypaimon.daft: OK")
配置 Catalog 连接
通过以下配置连接 DLF Catalog,本文示例统一引用变量名 CATALOG_OPTIONS:
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 | 固定为 |
uri | DLF Paimon REST 访问端点,详见 服务接入点与公网访问。VPC 访问支持 HTTP 和 HTTPS 协议,公网访问必须使用 HTTPS 协议。 |
warehouse | DLF Catalog 名称。 |
token.provider | AccessKey 认证模式下填 |
dlf.region | DLF 所属地域 ID(如 |
dlf.access-key-id / dlf.access-key-secret | DLF AccessKey。 |
dlf.security-token | (可选)STS 安全令牌。 |
dlf.oss-endpoint | (可选)OSS 公网访问 endpoint。 |
步骤 1:准备素材数据
从 ModelScope 加载公开的自动驾驶道路场景数据集,并提取 JPG 原图及其尺寸信息。
import io
from PIL import Image
from modelscope.msdatasets import MsDataset
ds = MsDataset.load("modelscope/image_object_detection_auto_dataset", split="test")
print(f"loaded {len(ds)} road-scene JPGs from ModelScope public dataset")
def fetch_jpgs(n: int = 24):
out = []
for i in range(min(n, len(ds))):
row = ds[i]
with open(row["Input Image:FILE"], "rb") as f:
jpg = f.read()
with Image.open(io.BytesIO(jpg)) as im:
w, h = im.size
out.append({"title": row["Title"], "bytes": jpg, "w": w, "h": h})
return out
raws = fetch_jpgs(24)每张图片约为 30-100KB JPG,均为真实公开数据。
步骤 2:创建 BLOB 分区表
channel 是典型查询过滤维度(如“取某个相机的某段帧”),将其用作 partition key 可让分区裁剪天然生效。
BLOB 表 Schema 限制:含 large_binary 列的表必须设置 row-tracking.enabled=true 与 data-evolution.enabled=true,否则 DLF 服务端将拒绝建表。
import pyarrow as pa
from pypaimon import CatalogFactory, Schema
catalog = CatalogFactory.create(CATALOG_OPTIONS)
catalog.create_database("default", True)
pa_schema = pa.schema([
pa.field("sample_token", pa.string(), nullable=False), # 帧唯一 ID
pa.field("scene_token", pa.string()), # 同一段场景的帧共享
pa.field("channel", pa.string()), # partition key
pa.field("timestamp", pa.timestamp("ms")),
pa.field("ego_speed_kmh", pa.float64()),
pa.field("weather", pa.string()),
pa.field("width", pa.int32()),
pa.field("height", pa.int32()),
pa.field("frame_jpg", pa.large_binary()), # 原图(BLOB)
pa.field("thumbnail", pa.large_binary()), # 缩略图(BLOB)
])
schema = Schema.from_pyarrow_schema(
pa_schema,
partition_keys=["channel"],
options={
"bucket": "-1",
"file.format": "parquet",
"row-tracking.enabled": "true", # BLOB 列必需
"data-evolution.enabled": "true", # BLOB 列必需
},
)
catalog.create_table("default.drive_frames", schema, True)步骤 3:写入多模态数据
合成 nuScenes 风格的元数据(channel、scene_token、timestamp、ego_speed_kmh、weather),bytes 字段填入真实图片。
import daft
from datetime import datetime, timedelta
from pypaimon.daft import write_paimon
CHANNELS = ["CAM_FRONT", "CAM_FRONT_LEFT", "CAM_FRONT_RIGHT"]
WEATHERS = ["clear", "cloudy", "rain"]
base_ts = datetime(2026, 5, 22, 9, 0, 0)
def make_thumbnail(jpg_bytes: bytes, max_side: int = 128) -> bytes:
with Image.open(io.BytesIO(jpg_bytes)) as im:
im.thumbnail((max_side, max_side))
buf = io.BytesIO()
im.convert("RGB").save(buf, format="JPEG", quality=70)
return buf.getvalue()
rows = []
for i, r in enumerate(raws):
rows.append({
"sample_token": r["title"],
"scene_token": f"scene_{i // 8:03d}",
"channel": CHANNELS[i % len(CHANNELS)],
"timestamp": base_ts + timedelta(milliseconds=i * 50),
"ego_speed_kmh": float(20 + (i * 3) % 80),
"weather": WEATHERS[(i // 6) % len(WEATHERS)],
"width": r["w"],
"height": r["h"],
"frame_jpg": r["bytes"],
"thumbnail": make_thumbnail(r["bytes"]),
})
arrow_tbl = pa.Table.from_pylist(rows, schema=pa_schema)
write_paimon(daft.from_arrow(arrow_tbl),
"default.drive_frames", CATALOG_OPTIONS, mode="append")步骤 4:按维度探索数据
4.1 分区裁剪
通过 channel 分区键过滤前向相机的帧,谓词自动下推至扫描层。
from pypaimon.daft import read_paimon
front = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.where(daft.col("channel") == "CAM_FRONT")
.sort("sample_token")
)
front.show()由于 channel 是分区键,谓词将自动下推,仅扫描该 channel 下的文件。
4.2 谓词 + 投影下推
按车速过滤高速帧,并仅投影所需的元数据列。
fast = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.where(daft.col("ego_speed_kmh") >= 60.0)
.select("sample_token", "channel", "ego_speed_kmh", "weather")
.sort("ego_speed_kmh", desc=True)
)
fast.show()ego_speed_kmh 上的 filter 会被下推到 Paimon 扫描层;由于 BLOB 列 frame_jpg 与 thumbnail 未出现在 SELECT 中,扫描会完全跳过这些列,全程不会触及 BLOB 字节。
4.3 按 BLOB 大小排序
large_binary 列被 pypaimon.daft 自动映射为 Daft 的 daft.File 引用类型,只暴露 path、offset、length,不物化字节,可直接按 length 排序。
df = read_paimon("default.drive_frames", CATALOG_OPTIONS)
all_rows = df.to_pydict()
top3 = sorted(
zip(all_rows["sample_token"], all_rows["channel"], all_rows["frame_jpg"]),
key=lambda t: t[2].length,
reverse=True,
)[:3]
for tk, ch, ref in top3:
print(f"{tk} ({ch}): {ref.length} bytes, path={ref.path}")实测输出如下:
d0c3fa90-d193ec71 (CAM_FRONT): 81852 bytes, path=oss://.../*.blob
d4316313-6a8d56d2 (CAM_FRONT_RIGHT): 81821 bytes, path=oss://.../*.blob
d4eb8adc-e95a5937 (CAM_FRONT_RIGHT): 80535 bytes, path=oss://.../*.blob步骤 5:按需读取 BLOB 字节
daft.File 仅提供引用信息,真正消费字节时需要根据 path、offset、length 从 OSS 拉取。
def fetch_blob(file_ref: daft.File, table) -> bytes:
"""读 BLOB 引用对应的字节。table 是 catalog.get_table(...) 拿到的表对象。"""
with table.file_io.new_input_stream(file_ref.path) as stream:
stream.seek(file_ref.offset)
return stream.read(file_ref.length)
# 用法
table = catalog.get_table("default.drive_frames")
one = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.where(daft.col("channel") == "CAM_FRONT")
.limit(1).to_pydict()
)
ref = one["frame_jpg"][0]
raw = fetch_blob(ref, table)
print(f"fetched {len(raw)} bytes")
with Image.open(io.BytesIO(raw)) as im:
print(f"decoded: {im.format} {im.size} {im.mode}")
# -> decoded: JPEG (1280, 720) RGB步骤 6:打 tag 与时间旅行
将“训练集 v1”打 tag 固化,便于后续按 tag 回溯到该时间点的数据。
tbl = catalog.get_table("default.drive_frames")
snaps = sorted(tbl.snapshot_manager().list_snapshots(), key=lambda s: s.id)
print("snapshots:", [s.id for s in snaps])
tbl.create_tag("training_v1", snaps[0].id)
# 回到那个时点
df_v1 = read_paimon("default.drive_frames", CATALOG_OPTIONS, tag_name="training_v1")
print("training_v1 frames:", df_v1.count_rows())也支持使用 read_paimon(..., snapshot_id=123) 按 snapshot id 进行时间旅行。
步骤 7:接入 ML 预处理 UDF
7.1 表达式派生列
通过 Daft 表达式直接派生宽高比列,无需额外计算节点。
enriched = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.select("sample_token", "channel", "width", "height")
.with_column("aspect_ratio", daft.col("width") / daft.col("height"))
)
enriched.show()7.2 UDF + groupby
使用 UDF 将连续车速分桶为低/中/高速分类,再按分类聚合统计。
@daft.func(return_dtype=daft.DataType.string())
def speed_bucket(speed_col):
out = []
for v in speed_col.to_pylist():
if v is None: out.append("unknown")
elif v < 40: out.append("low")
elif v < 80: out.append("mid")
else: out.append("high")
return out
stats = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.select("sample_token", "ego_speed_kmh")
.with_column("speed_class", speed_bucket(daft.col("ego_speed_kmh")))
.groupby("speed_class")
.count("sample_token")
.sort("speed_class")
)
stats.show()输出如下:
high: 4 frames
low: 7 frames
mid: 13 frames7.3 PIL 图像预处理 UDF(fetch + decode + resize + 抽取特征)
将第五步的 fetch_blob 封装进 UDF,按行完成 PIL 解码、resize 与亮度计算,构成典型感知模型预处理的雏形。
@daft.func(return_dtype=daft.DataType.float64())
def avg_brightness(file_col):
out = []
for ref in file_col.to_pylist():
if ref is None:
out.append(None); continue
raw = fetch_blob(ref, table) # 走 pypaimon FileIO
with Image.open(io.BytesIO(raw)) as im:
small = im.convert("L").resize((64, 64))
px = list(small.getdata())
out.append(sum(px) / len(px))
return out
bri = (
read_paimon("default.drive_frames", CATALOG_OPTIONS)
.where(daft.col("channel") == "CAM_FRONT")
.limit(4)
.select("sample_token", "frame_jpg")
.with_column("brightness", avg_brightness(daft.col("frame_jpg")))
.select("sample_token", "brightness")
.sort("brightness", desc=True)
)
bri.show()实测输出如下(亮度取值范围 0-255,前两张为白天场景,后两张为偏暗场景):
cdbd1882-be82474a: 112.12
d0c3fa90-d193ec71: 111.34
cf0b73a9-3474b8ca: 34.71
cad180c4-553ceeb1: 33.93将 avg_brightness 替换为 run_detection_model(frame_jpg) 等模型调用,即可在 Paimon 表上接出完整的感知预处理与推理流水线。
注意事项
BLOB 表 Schema 限制:含
large_binary列的表必须设置row-tracking.enabled=true与data-evolution.enabled=true,否则 DLF 服务端将拒绝建表。分区键选择:分区键需与典型查询 filter 对齐(如按
channel、按capture_date等),分区裁剪才能真正生效。BLOB 按需消费:
daft.File是 zero-copy 引用,不会触发 OSS 读取;仅在调用table.file_io.new_input_stream(...)时才会下载实际数据。