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

Data Lake Formation:PyPaimon を使用した Data Lake Formation (DLF) マルチモーダルテーブルの読み取り

最終更新日:Aug 01, 2026

PyPaimon を使用して DLF マルチモーダルテーブルから BLOB を読み取る方法について説明します。ローカル読み取りと、Ray ベースの分散読み取りの両方を扱います。

前提条件

開始する前に、次の依存関係をインストールします:

pip install pyjindosdk
pip install pypaimon-1.5.dev20260727.tar.gz
説明
  • PyPaimon パッケージを取得するには、pypaimon-1.5.dev20260727.tar.gz をクリックします。

  • 本番環境では pyjindosdk のインストールを推奨します。pyjindosdk を介して OSS から読み取ることで、レガシー OSS パスへのフォールバックを回避でき、高い同時実行での読み取りにおいても安定したパフォーマンスを維持できます。

接続の初期化

マルチモーダルテーブルを読み取る前に、カタログに接続してテーブルオブジェクトを取得します。以下のすべての読み取りシナリオ (ローカル、Ray 分散、Ray 上の Daft) は、この 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>")

接続パラメーター

パラメーター

説明

<DATABASE_NAME>

Paimon データベースの名前です。

<DLF_ENDPOINT>

DLF エンドポイント (例: http://<region-id>-vpc.dlf.aliyuncs.com)。詳細については、「DLF エンドポイント」をご参照ください。

<CATALOG_NAME>

DLF カタログの名前です。

<REGION_ID>

リージョン ID。例:cn-hangzhou。

<OSS_ENDPOINT>

OSS エンドポイント (たとえば、oss-<region-id>-internal.aliyuncs.com) です。 詳細については、「OSS のリージョンとエンドポイント」をご参照ください。

<ACCESS_KEY_ID>

Alibaba Cloud の AccessKey ID です。

<ACCESS_KEY_SECRET>

Alibaba Cloud の AccessKey Secret です。

<SECURITY_TOKEN>

STS 一時的な認証情報のセキュリティトークンです。長期 AccessKey ペアを使用する場合は、このパラメーターを省略します。

<TABLE_NAME>

マルチモーダルテーブルの名前です。

ローカル読み取り

小規模データセットの一括読み取り

メモリに収まる 1 つのクリップ、または少数のクリップを読み取るには、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,
    )
)

大量データのストリーミング

多数のクリップを、すべての BLOB バイトを一度にメモリへ読み込まずに読み取るには、stream_blobs を使用してバッチ単位でデータを処理します。

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)

パラメーター

パラメーター

説明

parallelism

BLOB を読み取るスレッド数です。

batch_size

stream_blobs でのバッチあたりの行数。デフォルト: 1024。テーブルレベルの動的パラメーター read.batch-size を使用して設定することもできます。

トレーニング用の連続したフレーム範囲の読み取り

トレーニングのワークロードでは、同一クリップから連続したフレームウィンドウ (例: 16、32、64 フレーム) が必要になることがよくあります。次の例は、単一クリップ内の連続したフレーム範囲を読み取ります。

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 ベースのワークロードでは、分散読み取りに 2 段階のアプローチを使用します:

  1. scan().to_ray() を呼び出して、スカラー列と BLOB 記述子をワーカーに分散して読み取ります。

  2. table.map_with_blobs() を呼び出して、各 Ray ワーカー上で BLOB バイトを読み取り、ユーザー定義関数 (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: clip_id や frame_index などのスカラー列を含む pyarrow.Table。
    blobs: BLOB 列名をキーとする dict[str, list[bytes | None]]。
    小さな pyarrow.Table を返す必要があります。副作用のみの処理の場合は空のテーブルを返します。
    """
    camera_0 = blobs["camera_0"]
    # Ray のオブジェクトストアで大きなペイロードが実体化されるのを防ぐため、未加工の BLOB バイトを返さないでください。
    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 は遅延評価のため、result_ds が消費されるまで BLOB は読み取られません。
for _ in result_ds.iter_batches(batch_format="pyarrow"):
    pass

パラメーター

パラメーター

説明

to_ray()

Ray は、利用可能なリソースとデータサイズに基づいて、読み取りの同時実行数とブロック数を自動的に決定します。通常、この設定は不要です。

parallelism

各 Ray タスク内で BLOB バイトを読み取るスレッド数です。Ray ワーカー数ではありません。デフォルト:64。

batch_size

1 回の呼び出しで UDF に渡される行数です。デフォルト:1024。BLOB が大きい場合、またはワーカーのメモリが限られている場合は、128、256、512 に減らしてください。

Ray タスクのリトライ設定

大規模なコールドリードのシナリオでは、map_with_blobs() の Ray タスクのリトライを設定します。

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,
    },
)

ディスクリプタの容量も大きい場合は、同じ ray_remote_args を to_ray() に渡します。

Ray 上の Daft を使用した分散 BLOB 処理

Ray 上の Daft は、分散読み取りの代替パスを提供します。すでに 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, 27)

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 記述子列をバイナリバイトに読み込みます。
# max_concurrency は各 Daft バッチ UDF 内の読み取りの同時実行数を制御します。
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)
    ],
)

# 例:バイトサイズのみを集計します。本番環境では、デコード、前処理、推論、書き込みのロジックを連結します。
# Ray のオブジェクトストアに大きなペイロードが入るのを防ぐため、未加工の BLOB バイトを保持または書き込まないでください。
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 バッチ UDF 内の読み取り同時実行性を制御します。

from pypaimon.daft import open_blob

def consume_one_blob(file):
    with open_blob(file, catalog_options, table_identifier) as stream:
        data = stream.read()
        # デコード、前処理、推論、書き込み
        return len(data)