All Products
Search
Document Center

Data Lake Formation:Gunakan Daft untuk membaca dan menulis data Iceberg multimodal

Last Updated:Jun 12, 2026

Topik ini menjelaskan cara menyimpan byte gambar secara langsung dalam tabel Iceberg DLF dan menggunakan Daft untuk menghasilkan gambar mini, menghitung penyematan (embedding), serta melakukan pencarian kemiripan visual.

Ikhtisar penyimpanan multimodal

Topik ini menggunakan pendekatan penyimpanan inline: byte mentah gambar atau audio/video, gambar mini, dan vektor penyematan disimpan bersama metadata dalam satu tabel Iceberg yang sama. Daft mendekode byte gambar langsung dari tabel tanpa mengakses jalur object storage. Operasi menulis, membaca, dan pencarian kemiripan semuanya dilakukan pada satu tabel.

Data

Tipe kolom

Raw image bytes

binary

Thumbnail (opsional; untuk penelusuran ringan)

binary

Embedding vector

list<float>

Catatan

Jika ukuran file media individual besar, penyimpanan inline akan meningkatkan ukuran file Parquet. Pilih pendekatan penyimpanan berdasarkan skala data Anda.

Persiapan lingkungan

Instal dependensi

  1. Python 3.10 atau versi yang lebih baru.

  2. Unduh paket PyIceberg yang kompatibel dengan DLF pyiceberg-0.10.0.dev1.tar.gz dan letakkan di direktori target.

  3. Instal PyIceberg dan dependensinya.

    pip install pyiceberg-0.10.0.dev1.tar.gz "pyarrow>=19.0.0,<22.0.0" "boto3>=1.24.59"
  4. Instal Daft.

    pip install "daft>=0.7.15"
Penting

PyArrow harus menggunakan versi sebelum 22 (pyarrow>=19.0.0,<22.0.0). Versi PyArrow 22 ke atas menyertakan AWS SDK yang secara default menambahkan checksum streaming aws-chunked selama unggahan. Antarmuka OSS yang kompatibel dengan S3 di DLF tidak mendukung encoding ini, sehingga operasi menulis akan gagal dengan error aws-chunked encoding is not supported.

Konfigurasikan parameter

Siapkan informasi berikut:

Parameter

Deskripsi

${accessKeyId}

ID AccessKey Akun Alibaba Cloud Anda.

${accessKeySecret}

Rahasia AccessKey Akun Alibaba Cloud Anda.

${regionId}

ID Wilayah tempat DLF dideploy, misalnya ap-southeast-1. Untuk daftar ID wilayah, lihat Endpoints.

${catalogName}

Nama katalog di DLF (berkorespondensi dengan gudang Iceberg).

${database}

Nama database target (namespace Iceberg).

Penting

Lindungi kredensial AccessKey Anda. Jangan hardcode kredensial tersebut dalam kode atau commit ke repositori kode. Kami menyarankan agar Anda mengambilnya dari variabel lingkungan atau layanan manajemen rahasia.

Catatan

Layanan REST Iceberg DLF hanya dapat diakses dari dalam VPC. Jalankan kode dalam topik ini dari lingkungan VPC di wilayah yang sama dengan DLF (misalnya Instance ECS atau kluster EMR). Untuk titik akhir setiap wilayah, lihat Iceberg REST endpoints.

Hubungkan ke katalog DLF

Gunakan load_catalog dari PyIceberg untuk terhubung ke katalog DLF melalui protokol REST Iceberg.

from pyiceberg.catalog import load_catalog

REGION = "${regionId}"

catalog = load_catalog(
    "dlf",
    **{
        "type": "rest",
        "uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
        "warehouse": "${catalogName}",
        "rest.signing-name": "DlfNext",
        "rest.signing-region": REGION,
        "rest.sigv4-enabled": "true",
        "client.access-key-id": "${accessKeyId}",
        "client.secret-access-key": "${accessKeySecret}",
        "client.region": REGION,
        "s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com",
    },
)

Siapkan contoh gambar

Siapkan dua contoh gambar JPEG lokal bernama n01.jpg dan n02.jpg (format RGB) dan letakkan di direktori saat ini. Anda dapat mengunduh satu gambar masing-masing dari kelas n01440764 (tench) dan n02979186 (cassette player) dalam dataset publik Imagenette dan mengganti namanya.

Buat tabel multimodal

Gunakan PyIceberg untuk membuat tabel multimodal dengan kolom untuk byte gambar dan vektor penyematan.

from pyiceberg.schema import Schema
from pyiceberg.types import (
    LongType, StringType, BinaryType, ListType, FloatType, NestedField,
)
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC

schema = Schema(
    NestedField(1, "image_id", LongType(), required=True),
    NestedField(2, "filename", StringType()),
    NestedField(3, "label", StringType()),
    NestedField(4, "image", BinaryType()),
    NestedField(5, "thumbnail", BinaryType()),
    NestedField(6, "embedding",
                ListType(element_id=7, element_type=FloatType(), element_required=False)),
)
table = catalog.create_table(
    ("${database}", "image_catalog"),
    schema=schema,
    partition_spec=UNPARTITIONED_PARTITION_SPEC,
)
Catatan

Pernyataan ini secara default membuat tabel format-versi 2. DLF mendukung tabel format-versi 3, tetapi PyIceberg dan Daft belum mendukung penulisan ke tabel v3 atau tipe VARIANT. Jangan atur properties={"format-version": "3"} untuk tabel multimodal.

Tulis data gambar

Simpan byte gambar mentah di kolom image dan gunakan fungsi gambar Daft untuk menghasilkan gambar mini dari byte tersebut. Daft versi 0.7.15 ke atas mengonfigurasi OSS secara otomatis, sehingga Anda tidak perlu meneruskan io_config secara manual.

import daft
from daft.functions import image

df = daft.from_pydict({
    "image_id": [1, 2],
    "filename": ["n01.jpg", "n02.jpg"],
    "label": ["tench", "cassette_player"],
    "image": [open("n01.jpg", "rb").read(),
              open("n02.jpg", "rb").read()],
})

# Hasilkan thumbnail dari byte gambar mentah: decode -> ubah ukuran menjadi 32x32 -> encode ulang sebagai JPEG
df = df.with_column("thumbnail",
        image.encode_image(image.resize(image.decode_image(df["image"]), 32, 32), "JPEG"))

df.write_iceberg(table, mode="append")

Baca dan proses gambar

Setelah Anda membaca tabel, dekode byte inline di kolom image secara langsung — tidak diperlukan akses ke object storage.

from daft.functions import image

df = daft.read_iceberg(table)
df = df.with_column("img", image.decode_image(df["image"]))
df = df.with_column("thumb", image.resize(df["img"], 64, 64))
df.show()

Hitung penyematan dan lakukan pencarian kemiripan

  1. Gunakan @daft.func untuk mengonversi gambar menjadi vektor guna menghasilkan kolom penyematan list<float>. Contoh berikut menggunakan deskriptor hasil downsampling; Anda dapat menggantinya dengan CLIP, ResNet, atau model lainnya.

    import numpy as np
    import daft
    from daft.functions import image
    
    @daft.func(return_dtype=daft.DataType.list(daft.DataType.float32()))
    def embed(img) -> list:
        v = np.asarray(img).astype("float32").reshape(-1)
        return (v / (np.linalg.norm(v) or 1.0)).tolist()
    
    # Dekode gambar inline, ubah ukuran ke dimensi seragam, dan hitung penyematan (juga dapat dipertahankan saat menulis; lihat contoh lengkap)
    df = daft.read_iceberg(table)
    df = df.with_column("embedding", embed(image.resize(image.decode_image(df["image"]), 8, 8)))
  2. Setelah Anda memiliki kolom embedding, panggil collect() untuk mengumpulkan data secara lokal, hitung kemiripan kosinus terhadap vektor gambar kueri, dan ambil hasil Top-K untuk pencarian visual.

    # Kumpulkan penyematan yang dihitung pada langkah sebelumnya
    d = df.select("image_id", "embedding").collect().to_pydict()
    vectors = {i: np.asarray(v, "float32") for i, v in zip(d["image_id"], d["embedding"])}
    
    # Kemiripan kosinus
    cos = lambda a, b: float(a @ b / ((np.linalg.norm(a) * np.linalg.norm(b)) or 1))
    
    # Gunakan image_id=1 sebagai gambar kueri dan ambil tetangga terdekat Top-K
    query_id = 1
    results = sorted(
        ((i, cos(vectors[query_id], v)) for i, v in vectors.items() if i != query_id),
        key=lambda r: -r[1],
    )
    print(results)

Contoh lengkap

Contoh berikut menjelaskan alur kerja end-to-end: terhubung ke DLF, buat tabel, tulis byte gambar beserta gambar mini dan penyematan, telusuri data, lakukan pencarian visual, dekode gambar, dan bersihkan.

import uuid
import numpy as np
import daft
from daft.functions import image
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import (
    LongType, StringType, BinaryType, ListType, FloatType, NestedField)
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC

REGION, CATALOG, DB = "${regionId}", "${catalogName}", "${database}"

# 1) Hubungkan ke DLF
catalog = load_catalog("dlf", **{
    "type": "rest", "uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
    "warehouse": CATALOG, "rest.signing-name": "DlfNext",
    "rest.signing-region": REGION, "rest.sigv4-enabled": "true",
    "client.access-key-id": "${accessKeyId}",
    "client.secret-access-key": "${accessKeySecret}", "client.region": REGION,
    "s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com"})

# Siapkan byte contoh gambar (letakkan dua gambar JPEG di direktori saat ini terlebih dahulu)
samples = [("n01.jpg", "tench", open("n01.jpg", "rb").read()),
           ("n02.jpg", "cassette_player", open("n02.jpg", "rb").read())]

# 2) Buat tabel multimodal
name = f"image_catalog_{uuid.uuid4().hex[:8]}"
table = catalog.create_table((DB, name), Schema(
    NestedField(1, "image_id", LongType(), required=True),
    NestedField(2, "filename", StringType()),
    NestedField(3, "label", StringType()),
    NestedField(4, "image", BinaryType()),
    NestedField(5, "thumbnail", BinaryType()),
    NestedField(6, "embedding",
                ListType(element_id=7, element_type=FloatType(), element_required=False))),
    partition_spec=UNPARTITIONED_PARTITION_SPEC)

@daft.func(return_dtype=daft.DataType.list(daft.DataType.float32()))
def embed(img) -> list:
    v = np.asarray(img).astype("float32").reshape(-1)
    return (v / (np.linalg.norm(v) or 1.0)).tolist()

try:
    # 3) Tulis byte gambar dengan gambar mini dan penyematan turunan
    rows = {"image_id": [], "filename": [], "label": [], "image": []}
    for i, (fn, label, jpg) in enumerate(samples, 1):
        rows["image_id"].append(i)
        rows["filename"].append(fn)
        rows["label"].append(label)
        rows["image"].append(jpg)
    df = daft.from_pydict(rows)
    df = df.with_column("thumbnail",
            image.encode_image(image.resize(image.decode_image(df["image"]), 32, 32), "JPEG"))
    df = df.with_column("embedding",
            embed(image.resize(image.decode_image(df["image"]), 8, 8)))
    df.select("image_id", "filename", "label", "image", "thumbnail",
              "embedding").write_iceberg(table, mode="append")

    # 4) Telusuri data (projection pushdown, tanpa kolom piksel)
    table = catalog.load_table((DB, name))
    daft.read_iceberg(table).select("image_id", "label", "filename").sort("image_id").show()

    # 5) Pencarian visual: temukan tetangga kosinus untuk image_id=1
    d = daft.read_iceberg(table).select("image_id", "embedding").collect().to_pydict()
    M = {i: np.asarray(v, "float32") for i, v in zip(d["image_id"], d["embedding"])}
    cos = lambda a, b: float(a @ b / ((np.linalg.norm(a) * np.linalg.norm(b)) or 1))
    print(sorted(((i, cos(M[1], v)) for i, v in M.items() if i != 1), key=lambda r: -r[1]))

    # 6) Ambil dan dekode gambar
    one = daft.read_iceberg(table).where(daft.col("image_id") == 1)
    one = one.with_column("img", image.decode_image(one["image"]))
    print("decoded shape:", one.select("img").collect().to_pydict()["img"][0].shape)
finally:
    catalog.drop_table((DB, name))
Catatan

Penyematan dalam contoh ini adalah deskriptor hasil downsampling yang hanya digunakan untuk tujuan demonstrasi. Di lingkungan produksi, gantilah dengan CLIP, ResNet, atau model lainnya untuk inferensi.