All Products
Search
Document Center

Data Lake Formation:Baca tabel multimodal DLF menggunakan PyPaimon

Last Updated:Aug 01, 2026

Pelajari cara membaca BLOB dari tabel multimodal DLF dengan PyPaimon, mencakup pembacaan lokal maupun terdistribusi berbasis Ray.

Prasyarat

Instal dependensi berikut sebelum memulai:

pip install pyjindosdk
pip install pypaimon-1.5.dev20260727.tar.gz
Catatan
  • Untuk mendapatkan paket PyPaimon, klik pypaimon-1.5.dev20260727.tar.gz.

  • Kami merekomendasikan menginstal pyjindosdk di lingkungan produksi. Membaca dari OSS melalui pyjindosdk menghindari fallback ke path OSS lama, sehingga menjaga performa stabil dalam kondisi pembacaan konkurensi tinggi.

Inisialisasi koneksi

Sebelum membaca tabel multimodal, hubungkan ke katalog dan peroleh objek tabel. Semua skenario baca berikut — lokal, terdistribusi Ray, dan Daft pada Ray — menggunakan kembali objek table ini.

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>")

Parameter koneksi

Parameter

Deskripsi

<DATABASE_NAME>

Nama database Paimon.

<DLF_ENDPOINT>

Titik akhir DLF, contohnya http://<region-id>-vpc.dlf.aliyuncs.com. Untuk informasi lebih lanjut, lihat DLF Endpoint.

<CATALOG_NAME>

Nama katalog DLF.

<REGION_ID>

ID Wilayah, contohnya cn-hangzhou.

<OSS_ENDPOINT>

Titik akhir OSS, contohnya oss-<region-id>-internal.aliyuncs.com. Untuk informasi lebih lanjut, lihat Wilayah dan titik akhir OSS.

<ACCESS_KEY_ID>

ID AccessKey Alibaba Cloud Anda.

<ACCESS_KEY_SECRET>

Rahasia AccessKey Alibaba Cloud Anda.

<SECURITY_TOKEN>

Token keamanan dari kredensial temporary STS. Abaikan parameter ini jika Anda menggunakan Pasangan Kunci Akses jangka panjang.

<TABLE_NAME>

Nama tabel multimodal.

Pembacaan lokal

Baca set data kecil sekaligus

Untuk membaca satu klip atau beberapa klip yang muat dalam memori, gunakan read_blobs untuk memuat semua data sekaligus.

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

Alirkan volume data besar

Untuk membaca banyak klip tanpa memuat seluruh byte BLOB ke dalam memori sekaligus, gunakan stream_blobs untuk mengonsumsi data dalam batch.

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)

Parameter

Parameter

Deskripsi

parallelism

Jumlah thread untuk membaca BLOB.

batch_size

Jumlah baris per batch dalam stream_blobs. Nilai default: 1024. Anda juga dapat mengatur nilai ini melalui parameter dinamis tingkat tabel read.batch-size.

Baca rentang frame berurutan untuk pelatihan

Beban kerja pelatihan sering kali memerlukan jendela frame berurutan — misalnya, 16, 32, atau 64 frame — dari klip yang sama. Contoh berikut membaca rentang frame berurutan dalam satu klip.

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)

Pemrosesan BLOB terdistribusi dengan Ray

Untuk beban kerja berbasis Ray, gunakan pendekatan dua langkah untuk pembacaan terdistribusi:

  1. Panggil scan().to_ray() untuk membaca kolom skalar dan deskriptor BLOB di seluruh worker.

  2. Panggil table.map_with_blobs() untuk membaca byte BLOB pada setiap worker Ray dan meneruskannya ke user-defined function (UDF).

Penggunaan dasar

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 yang berisi kolom skalar seperti clip_id dan frame_index.
    blobs: dict[str, list[bytes | None]], diindeks berdasarkan nama kolom blob.
    Harus mengembalikan pyarrow.Table kecil. Kembalikan tabel kosong jika hanya untuk efek samping.
    """
    camera_0 = blobs["camera_0"]
    # Hindari mengembalikan byte BLOB mentah untuk mencegah materialisasi muatan besar di object store Ray.
    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,
)

# Jalankan eksekusi. Ray Dataset bersifat lazy — BLOB tidak dibaca hingga result_ds dikonsumsi.
for _ in result_ds.iter_batches(batch_format="pyarrow"):
    pass

Parameter

Parameter

Deskripsi

to_ray()

Ray secara otomatis menentukan konkurensi baca dan jumlah blok berdasarkan sumber daya yang tersedia dan ukuran data. Biasanya Anda tidak perlu mengonfigurasi ini.

parallelism

Jumlah thread dalam setiap task Ray untuk membaca byte BLOB. Ini bukan jumlah worker Ray. Nilai default: 64.

batch_size

Jumlah baris yang diteruskan ke UDF per pemanggilan. Nilai default: 1024. Jika BLOB berukuran besar atau memori worker terbatas, kurangi nilainya menjadi 128, 256, atau 512.

Konfigurasi retry task Ray

Untuk skenario cold-read skala besar, konfigurasikan retry task Ray untuk map_with_blobs().

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

Jika volume deskriptor juga besar, teruskan argumen ray_remote_args yang sama ke to_ray().

Pemrosesan BLOB terdistribusi dengan Daft pada Ray

Daft pada Ray menyediakan jalur baca terdistribusi alternatif, cocok digunakan jika Anda sudah memiliki pipeline berbasis Daft.

Penggunaan dasar

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 membaca kolom deskriptor File/Blob Daft menjadi byte biner.
# max_concurrency mengontrol konkurensi baca dalam setiap UDF batch Daft.
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)
    ],
)

# Contoh: agregasi hanya ukuran byte. Di produksi, rangkai logika decode, preprocess, inferensi, atau write.
# Hindari menyimpan atau menulis byte BLOB mentah untuk mencegah muatan besar masuk ke object store Ray.
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()

Buka aliran blob individual dengan open_blob

Untuk membuka dan memproses BLOB secara individual di dalam UDF Daft kustom, gunakan open_blob(). Parameter max_concurrency mengontrol konkurensi baca dalam setiap UDF batch Daft.

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 / inferensi / write
        return len(data)