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

Data Lake Formation:Daft を使用したマルチモーダル Iceberg データの読み取りと書き込み

最終更新日:Jul 11, 2026

このトピックでは、DLF Iceberg テーブルに画像バイトを直接格納し、Daft を使用してサムネイル生成、埋め込みベクトルの計算、視覚的な類似検索を行う方法について説明します。

マルチモーダルストレージの概要

このトピックでは、インライン格納アプローチを使用します。生の画像または音声/動画のバイト、サムネイル、埋め込みベクトルを、メタデータと同じ Iceberg テーブルに格納します。Daft は、オブジェクトストレージのパスにアクセスせずに、テーブルから画像バイトを直接デコードします。書き込み、読み取り、類似検索はすべて 1 つのテーブルに対して行われます。

データ

列タイプ

生の画像バイト

binary

サムネイル (任意:軽量なブラウジング用)

binary

埋め込みベクトル

list<float>

説明

個々のメディアファイルが大きい場合、インライン格納により Parquet ファイルサイズが増加します。データの規模に応じて格納方式を選択してください。

環境のセットアップ

依存関係のインストール

  1. Python 3.10 以降。

  2. DLF 互換の PyIceberg パッケージ (pyiceberg-dlf) をインストールします。

    python3 -m venv venv
    source venv/bin/activate
    pip install -U pip
    # pyiceberg をアンインストール (pyiceberg-dlf と共存できないため)
    pip uninstall -y pyiceberg
    # rest-sigv4 が必要 (REST sigv4 署名用に boto3 をインストール)
    pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"
  3. Daft をインストールします。

    pip install "daft>=0.7.17"

パラメータの設定

以下の情報を準備してください。

パラメータ

説明

${accessKeyId}

AccessKey ID。

${accessKeySecret}

AccessKey secret。

${regionId}

DLF がデプロイされているリージョン ID です。例: ap-southeast-1。リージョン ID の一覧については、「エンドポイントとパブリックネットワークアクセス」をご参照ください。

${catalogName}

DLF のカタログ名 (Iceberg ウェアハウスに対応)。

${database}

対象データベースの名前 (Iceberg 名前空間)。

重要

AccessKey 認証情報を保護してください。コードにハードコードしたり、コードリポジトリにコミットしたりしないでください。環境変数またはシークレット管理サービスから取得することを推奨します。

DLF カタログへの接続

PyIceberg の load_catalog を使用して、Iceberg REST プロトコル経由で DLF カタログに接続します。

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

サンプル画像の準備

n01.jpg と n02.jpg という名前のローカル JPEG サンプル画像 (RGB 形式) を 2 つ用意し、カレントディレクトリに配置してください。公開 Imagenette データセットの n01440764 (tench) クラスと n02979186 (cassette player) クラスからそれぞれ 1 枚ずつ画像をダウンロードし、リネームして使用できます。

マルチモーダルテーブルの作成

PyIceberg を使用して、画像バイトと埋め込みベクトルの列を含むマルチモーダルテーブルを作成します。

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

このステートメントは、デフォルトでフォーマットバージョン 2 のテーブルを作成します。DLF はフォーマットバージョン 3 のテーブルをサポートしていますが、PyIceberg と Daft はまだ v3 テーブルへの書き込みや VARIANT 型をサポートしていません。マルチモーダルテーブルには properties={"format-version": "3"} を設定しないでください。

画像データの書き込み

生の画像バイトを image 列に保存し、Daft の画像関数を使用してそれらのバイトからサムネイルを生成します。Daft 0.7.15 以降では OSS が自動的に設定されるため、io_config を手動で渡す必要はありません。

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

# 生の画像バイトからサムネイルを派生:デコード → 32x32 にリサイズ → 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")

画像の読み取りと処理

表を確認した後、image 列のインラインバイトを直接デコードします。オブジェクトストレージへのアクセスは必要ありません。

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

埋め込みの計算と類似検索

  1. @daft.func を使用して画像をベクトルに変換し、list<float> の埋め込み列を生成します。次の例では、ダウンサンプリングされた記述子を使用します。これを CLIP、ResNet、または別のモデルで置き換えることができます。

    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()
    
    # インライン画像をデコードし、サイズを統一して埋め込みを計算 (書き込み時に永続化することも可能。完全な例を参照)
    df = daft.read_iceberg(table)
    df = df.with_column("embedding", embed(image.resize(image.decode_image(df["image"]), 8, 8)))
  2. embedding 列を取得したら、collect() を呼び出してデータをローカルに収集し、クエリ画像ベクトルとのコサイン類似度を計算し、画像検索の上位 K 件の結果を取得します。

    # 前のステップで計算した埋め込みを収集
    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"])}
    
    # コサイン類似度
    cos = lambda a, b: float(a @ b / ((np.linalg.norm(a) * np.linalg.norm(b)) or 1))
    
    # image_id=1 をクエリ画像として使用し、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)

完全な例

次の例では、エンドツーエンドのワークフローを順に実行します。DLF への接続、テーブルの作成、サムネイルと埋め込みを含む画像バイトの書き込み、データの参照、視覚検索の実行、画像のデコード、クリーンアップ。

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

# サンプル画像バイトを準備 (事前にカレントディレクトリに 2 つの JPEG 画像を配置)
samples = [("n01.jpg", "tench", open("n01.jpg", "rb").read()),
           ("n02.jpg", "cassette_player", open("n02.jpg", "rb").read())]

# 2) マルチモーダルテーブルを作成
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) サムネイルと埋め込みを派生させて画像バイトを書き込み
    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) データを参照 (射影プッシュダウン、ピクセル列なし)
    table = catalog.load_table((DB, name))
    daft.read_iceberg(table).select("image_id", "label", "filename").sort("image_id").show()

    # 5) 視覚検索: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) 画像を取得してデコード
    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))
説明

この例の埋め込みは、デモ用のダウンサンプリング記述子にすぎません。本番環境では、推論向けに CLIP、ResNet などのモデルに置き換えてください。