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 |
|
|
Thumbnail (opsional; untuk penelusuran ringan) |
|
|
Embedding vector |
|
Jika ukuran file media individual besar, penyimpanan inline akan meningkatkan ukuran file Parquet. Pilih pendekatan penyimpanan berdasarkan skala data Anda.
Persiapan lingkungan
Instal dependensi
-
Python 3.10 atau versi yang lebih baru.
-
Unduh paket PyIceberg yang kompatibel dengan DLF pyiceberg-0.10.0.dev1.tar.gz dan letakkan di direktori target.
-
Instal PyIceberg dan dependensinya.
pip install pyiceberg-0.10.0.dev1.tar.gz "pyarrow>=19.0.0,<22.0.0" "boto3>=1.24.59" -
Instal Daft.
pip install "daft>=0.7.15"
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 |
|
|
ID AccessKey Akun Alibaba Cloud Anda. |
|
|
Rahasia AccessKey Akun Alibaba Cloud Anda. |
|
|
ID Wilayah tempat DLF dideploy, misalnya |
|
|
Nama katalog di DLF (berkorespondensi dengan gudang Iceberg). |
|
|
Nama database target (namespace Iceberg). |
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.
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,
)
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
-
Gunakan
@daft.funcuntuk mengonversi gambar menjadi vektor guna menghasilkan kolom penyematanlist<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))) -
Setelah Anda memiliki kolom
embedding, panggilcollect()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))
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.