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
-
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 |
|
|
Nama database Paimon. |
|
|
Titik akhir DLF, contohnya |
|
|
Nama katalog DLF. |
|
|
ID Wilayah, contohnya |
|
|
Titik akhir OSS, contohnya |
|
|
ID AccessKey Alibaba Cloud Anda. |
|
|
Rahasia AccessKey Alibaba Cloud Anda. |
|
|
Token keamanan dari kredensial temporary STS. Abaikan parameter ini jika Anda menggunakan Pasangan Kunci Akses jangka panjang. |
|
|
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 |
|
|
Jumlah thread untuk membaca BLOB. |
|
|
Jumlah baris per batch dalam |
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:
-
Panggil
scan().to_ray()untuk membaca kolom skalar dan deskriptor BLOB di seluruh worker. -
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 |
|
|
Ray secara otomatis menentukan konkurensi baca dan jumlah blok berdasarkan sumber daya yang tersedia dan ukuran data. Biasanya Anda tidak perlu mengonfigurasi ini. |
|
|
Jumlah thread dalam setiap task Ray untuk membaca byte BLOB. Ini bukan jumlah worker Ray. Nilai default: 64. |
|
|
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)