全部產品
Search
文件中心

Data Lake Formation:基於Daft和DLF的自動駕駛多模態資料湖實踐

更新時間:Jul 10, 2026

示範如何用 Daft + DLF Paimon 表管理自動駕駛多網路攝影機幀資料:原圖與縮圖以 BLOB 形式入湖、按 channel/車速/天氣快速過濾、按需拉取大位元組、為訓練集打 tag 回溯,並接入 PIL/模型預先處理 UDF。

情境

智能駕駛車隊從多網路攝影機採集幀序列入湖,下遊用於訓練集篩選、迴歸集抽樣、感知模型批量預先處理。典型需求如下:

  • 按 channel(相機通道)、時間窗、自車狀態等維度快速定位幀。

  • BLOB 形式的映像不能每次都全量物化。

  • 訓練集快照(snapshot)可回溯、可打 tag(tag 即 snapshot 的命名引用)。

各組件分工如下:

組件

提供能力

Paimon

儲存、Schema 演化、時間旅行

Daft

lazy DataFrame、最佳化器、BLOB-as-daft.File 多模態語義

DLF

Catalog 服務、OSS 憑證管理

前提條件

安裝依賴

  1. 下載 pypaimon 離線安裝包pypaimon-1.5.dev20260704.tar.gz

  2. 安裝 Python 依賴

    在離線包所在目錄執行以下命令,一次性安裝通用依賴與本情境額外依賴:

    pip install \
        daft \
        'pyarrow>=16,<=21' \
        pandas \
        requests \
        pypaimon-1.5.dev20260608.tar.gz \
        modelscope \
        datasets \
        pillow \
        oss2
  3. 驗證環境

    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

固定為 rest

uri

DLF Paimon REST 訪問端點,詳見 服務存取點與公網訪問。VPC 訪問支援 HTTP 和 HTTPS 協議,公網訪問必須使用 HTTPS 協議。

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。

步驟 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=truedata-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 風格的中繼資料(channelscene_tokentimestampego_speed_kmhweather),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_jpgthumbnail 未出現在 SELECT 中,掃描會完全跳過這些列,全程不會觸及 BLOB 位元組。

4.3 按 BLOB 大小排序

large_binary 列被 pypaimon.daft 自動對應為 Daft 的 daft.File 參考型別,只暴露 pathoffsetlength,不物化位元組,可直接按 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 僅提供引用資訊,真正消費位元組時需要根據 pathoffsetlength 從 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 frames

7.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=truedata-evolution.enabled=true,否則 DLF 服務端將拒絕建表。

  • 分區鍵選擇:分區鍵需與典型查詢 filter 對齊(如按 channel、按 capture_date 等),分區裁剪才能真正生效。

  • BLOB 按需消費:daft.File 是 zero-copy 引用,不會觸發 OSS 讀取;僅在調用 table.file_io.new_input_stream(...) 時才會下載實際資料。