すべてのプロダクト
Search
ドキュメントセンター

Data Lake Formation:Daft を用いた DLF Paimon でのマルチモーダル自動運転データの管理

最終更新日:Jun 06, 2026

Data Lake Formation (DLF) Paimon テーブルで Daft を使用して、マルチカメラの自動運転フレームを管理します。ソース画像とサムネイルを BLOB としてインジェストし、チャネル、自車速、天候などでフィルタリングし、バイトデータをオンデマンドで取得し、タイムトラベル用にトレーニングセットのスナップショットをタギングし、PIL やモデルの前処理をユーザー定義関数 (UDF) として組み込みます。

シナリオ

自動運転のフリートは、マルチカメラのフレームシーケンスをデータレイクにインジェストし、下流のトレーニングセット選択、回帰サンプリング、バッチでの知覚前処理に供給します。一般的な要件は次のとおりです。

  • channel (カメラチャネル)、時間枠、自車両の状態などのディメンションでフレームを迅速に特定する。

  • 読み取りごとに画像 BLOB を完全にマテリアライズすることを回避する。

  • 後で名前で参照できるように、トレーニングセットのスナップショットをタギングする。

コンポーネントの責務:

コンポーネント

機能

Paimon

ストレージ、スキーマエボリューション、タイムトラベル

Daft

遅延 DataFrame、オプティマイザ、daft.File としての BLOB のマルチモーダルセマンティクス

DLF

カタログサービス、OSS の認証情報管理

前提条件

依存関係のインストール

  1. pypaimon オフラインパッケージのダウンロード

    お使いのブラウザで、以下のリンクから pypaimon のオフラインインストールパッケージをダウンロードします:[TODO: 永続的なオフラインパッケージの URL は保留中]。

  2. Python 依存関係のインストール

    オフラインパッケージを格納しているディレクトリで、次のコマンドを実行して、共通の依存関係とこのチュートリアルで必要な追加機能をインストールします。

    pip install \
        daft \
        'pyarrow>=16,<=21' \
        pandas \
        requests \
        pypaimon-1.5.dev20260526.tar.gz \
        modelscope \
        datasets \
        pillow \
        oss2
  3. 環境の確認

    import daft
    import pypaimon
    from pypaimon.daft import read_paimon, write_paimon
    print("daft:", daft.__version__, "pypaimon:", pypaimon.__version__)

カタログ接続の設定

DLF カタログに接続するには、以下の設定を使用します。変数名 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 カタログ名。

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 インターネットエンドポイント。

ステップ 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)

これらは実世界の JPG で、それぞれ約 30~100 KB です。

ステップ 2:パーティション分割された BLOB テーブルの作成

channel は典型的なフィルターディメンション (例:「1 つのカメラから特定の時間枠のフレームを取得する」) であるため、これをパーティションキーにすることで、パーティションプルーニングが自動的に機能します。

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()),                      # 1 つのシーンのフレームで共有
    pa.field("channel",       pa.string()),                      # パーティションキー
    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 フィールドに実際の JPG バイトを格納します。

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 はパーティションキーであるため、述語はプッシュダウンされ、そのチャネルのファイルのみがスキャンされます。

4.2 述語とプロジェクションのプッシュダウン

高速なフレームをフィルタリングし、必要なメタデータ列のみを射影 (project) します。

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 に対するフィルターは Paimon のスキャンレイヤーにプッシュダウンされます。BLOB 列である frame_jpgthumbnail は SELECT リストに含まれていないため、スキャンはこれらの列を完全にスキップし、BLOB のバイトデータに触れることはありません。

4.3 BLOB サイズによるソート

pypaimon.daftlarge_binary 列を 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:タギングとタイムトラベル

現在のスナップショットを training_v1 としてタギングし、後でこの状態に戻れるようにします。

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) を使用して、スナップショット ID によるタイムトラベルも可能です。

ステップ 7:ML 前処理を UDF として組み込む

7.1 式からの派生列

Daft の式を直接使用してアスペクト比の列を派生させます。UDF は不要です。

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 (取得 + デコード + リサイズ + 特徴抽出)

ステップ 5 の fetch_blob を UDF でラップし、PIL でデコード、リサイズ、各行の輝度を計算します。これは、知覚前処理パイプラインの典型的なパターンです。

@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。最初の 2 つは日中のシーン、最後の 2 つは薄暗いシーンです):

cdbd1882-be82474a: 112.12
d0c3fa90-d193ec71: 111.34
cf0b73a9-3474b8ca:  34.71
cad180c4-553ceeb1:  33.93

avg_brightnessrun_detection_model(frame_jpg) や別のモデル呼び出しに置き換えることで、Paimon テーブルで完全な知覚前処理と推論パイプラインを構築できます。

注意事項

  • BLOB テーブルのスキーマ制約:large_binary 列を含むテーブルは、row-tracking.enabled=true および data-evolution.enabled=true を設定する必要があります。設定しない場合、DLF サービスはテーブル作成リクエストを拒否します。

  • パーティションキーの選択:パーティションプルーニングが実際に機能するように、パーティションキーを典型的なフィルター (例:channelcapture_date) に合わせます。

  • BLOB のオンデマンド消費:daft.File はゼロコピーの参照であり、OSS の読み取りをトリガーしません。実際のダウンロードは table.file_io.new_input_stream(...) を呼び出したときにのみ発生します。