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

Data Lake Formation:Python を使用した DLF Lance テーブルの操作

最終更新日:May 13, 2026

このガイドでは、lance-dlf を使用して、DLF カタログに接続し、Lance テーブルを作成し、データを書き込み、結果を検証する方法について説明します。

説明

Daft を使用して DLF の Lance テーブルを操作するには、「Daft を使用して DLF Lance テーブルを操作する」をご参照ください。

機能概要

lance-dlf は、以下のコア機能を提供します:

  • DLF カタログへの接続

  • DLF データベースを Lance 名前空間にマッピング (論理マッピングレイヤー)

  • type=lance-table の DLF テーブルのみを公開します

  • DLF の load_table_token API 経由で一時的な OSS アクセス資格情報を取得します

  • OSS の一時的な資格情報を PyLance 用の storage_options に変換します

PyLance が実際のデータの読み取り/書き込み操作を処理します。

# データの書き込み
lance.write_dataset(table, location, storage_options=storage_options)

# データの読み取り
lance.dataset(location, storage_options=storage_options)

クイックスタート

lance-dlf のインストール

PyPI から lance-dlf をインストールします:

python3 -m pip install lance-dlf

カタログ接続の設定

CONFIG = {
    "uri": "http://<dlf-endpoint>", # パブリックアクセスの場合は HTTPS プロトコルを使用
    "warehouse": "<warehouse>",
    "token.provider": "dlf",
    "dlf.region": "<region>",
    "dlf.access-key-id": "<access-key-id>",
    "dlf.access-key-secret": "<access-key-secret>",
    "dlf.oss-endpoint": "<oss-endpoint>",
}

設定パラメーター:

パラメーター

説明

uri

DLF Paimon REST エンドポイント。パブリックアクセスには HTTPS プロトコルを使用します

ウェアハウス

DLF カタログ名

token.provider

dlf を使用した AccessKey 認証

dlf.region

DLF リージョン ID、例: ap-southeast-1

dlf.access-key-id

DLF アクセス用の AccessKey ID

dlf.access-key-secret

DLF アクセス用の AccessKey Secret

dlf.security-token

(オプション) STS シナリオ用のセキュリティトークン

dlf.oss-endpoint

(任意) OSS パブリック エンドポイント (例: oss-ap-southeast-1.aliyuncs.com)。 デフォルトで無効になっているパブリック DLF アクセスに必要です。 DLF カタログのパブリックアクセスを有効にするには、「DLF のパブリックネットワーク接続が利用可能になりました」をご参照ください。

重要

AccessKey ID と AccessKey Secret は、Alibaba Cloud リソースにアクセスするための重要な認証情報です。安全に保管し、実際の AccessKey を Git リポジトリにコミットしないでください。認証情報は、以下のいずれかの場所から読み込むようにしてください。

  • 環境変数

  • キー管理システム

  • ランタイム構成

DLF カタログへの接続

lance_dlf モジュールをインポートすると、dlf 名前空間の実装が自動的に登録されます:

import lance_namespace
import lance_dlf  # noqa: F401

# DLF カタログに接続
ns = lance_namespace.connect("dlf", CONFIG)

# 接続を確認
print(ns.namespace_id())

namespace_id() メソッドは、DLF エンドポイントとウェアハウスを含む、接続されているカタログに関する情報を返します。

基本操作

新しい Lance テーブルの作成とデータの書き込み

テーブルが存在しない場合、テーブルを作成してデータを挿入します。メカニズム:

  1. ns.create_table() を使用してテーブルを作成します

  2. DLF から Lance テーブルのストレージの場所(location)を取得する

  3. lance-dlf 経由で一時的な OSS アクセス資格情報を取得する

  4. PyLance を使用して Arrow データを IPC バイトに変換し、指定された場所にデータを書き込みます (「リファレンス」をご参照ください)

import lance
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest

DATABASE = "default"
TABLE = "test_lance_create_001"

table_id = [DATABASE, TABLE]

# テストデータを構築
data = pa.table({
    "f0": pa.array([101, 102, 103], type=pa.int64()),
    "f1": pa.array(["create-a", "create-b", "create-c"], type=pa.string()),
})

# Arrow テーブルを IPC バイトに変換
def arrow_table_to_ipc_bytes(table):
    sink = pa.BufferOutputStream()
    writer = pa.ipc.new_stream(sink, table.schema)
    writer.write_table(table)
    writer.close()
    return sink.getvalue().to_pybytes()

# テーブルを作成してデータを書き込む
create_response = ns.create_table(
    CreateTableRequest(id=table_id),
    arrow_table_to_ipc_bytes(data),
)

print(create_response.location)
print(create_response.storage_options.keys())

# 読み取りと検証
desc = ns.describe_table(DescribeTableRequest(id=table_id))
dataset = lance.dataset(desc.location, storage_options=desc.storage_options)
result = dataset.to_table()
print(result)

期待される出力:

pyarrow.Table
f0: int64
f1: string
----
f0: [[101,102,103]]
f1: [["create-a","create-b","create-c"]]

既存の空のテーブルへの書き込み

DLF に空の type=lance-table テーブルが存在する場合、describe_table を使用してテーブルのストレージロケーションとアクセス認証情報を取得し、PyLance を使用してデータを書き込みます。

import lance
import pyarrow as pa
from lance_namespace import DescribeTableRequest

DATABASE = "default"
TABLE = "test_lance_table"

table_id = [DATABASE, TABLE]

# テーブルの詳細を取得
desc = ns.describe_table(DescribeTableRequest(id=table_id))

# テストデータを構築
data = pa.table({
    "f0": pa.array([1, 2, 3], type=pa.int64()),
    "f1": pa.array(["value-1", "value-2", "value-3"], type=pa.string()),
})

# データを書き込む
lance.write_dataset(
    data,
    desc.location,
    mode="overwrite",
    storage_options=desc.storage_options,
)

# 読み取りと検証
dataset = lance.dataset(desc.location, storage_options=desc.storage_options)
print(dataset.to_table())

書き込みモード:

モード

説明

上書き

既存のデータを上書きします。空のテーブルやテストテーブルの初期化に便利です。

append

データを追記します。スキーマに互換性がある必要があります。

警告

overwrite モードは、既存のすべての Lance データセットのデータを上書きします。このモードを使用する際は注意してください。

名前空間とテーブルの表示

from lance_namespace import (
    DescribeNamespaceRequest,
    DescribeTableRequest,
    ListNamespacesRequest,
    ListTablesRequest,
)

DATABASE = "<database>"
TABLE = "<table>"

# すべての名前空間を一覧表示
namespaces = ns.list_namespaces(ListNamespacesRequest(id=[]))
print(namespaces)

# 名前空間の詳細を取得
namespace = ns.describe_namespace(DescribeNamespaceRequest(id=[DATABASE]))
print(namespace)

# データベース内のすべてのテーブルを一覧表示
tables = ns.list_tables(ListTablesRequest(id=[DATABASE]))
print(tables)

# テーブルの詳細を取得
table = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
print(table.location)
print(table.properties)
print(table.storage_options)

describe_table の戻り値フィールド:

フィールド

説明

location

Lance データセットの物理ストレージパス (通常は oss://bucket/path)

プロパティ

DLF テーブルスキーマのオプションでは、type フィールドは lance-table である必要があります

storage_options

PyLance が OSS を読み書きするための一時的なアクセス認証情報

この例では、DLF への接続、テーブルの作成、データの書き込み、テーブルの一覧表示、結果の検証という完全なワークフローを示します。

from datetime import datetime

import lance
import lance_namespace
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest, ListTablesRequest

import lance_dlf  # noqa: F401


CONFIG = {
    "uri": "http://<dlf-endpoint>",
    "warehouse": "<warehouse>",
    "token.provider": "dlf",
    "dlf.region": "<region>",
    "dlf.access-key-id": "<access-key-id>",
    "dlf.access-key-secret": "<access-key-secret>",
    "dlf.oss-endpoint": "<oss-endpoint>",
}

DATABASE = "default"


def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
    """PyArrow テーブルを IPC バイトストリームに変換"""
    sink = pa.BufferOutputStream()
    with pa.ipc.new_stream(sink, table.schema) as writer:
        writer.write_table(table)
    return sink.getvalue().to_pybytes()


def main():
    # DLF カタログに接続
    ns = lance_namespace.connect("dlf", CONFIG)

    # 一意のテーブル名を生成
    table_name = "test_lance_create_" + datetime.now().strftime("%Y%m%d_%H%M%S")
    table_id = [DATABASE, table_name]

    # テストデータを構築
    data = pa.table({
        "f0": pa.array([101, 102, 103], type=pa.int64()),
        "f1": pa.array(["create-a", "create-b", "create-c"], type=pa.string()),
    })

    # テーブルを作成してデータを書き込む
    create_response = ns.create_table(
        CreateTableRequest(id=table_id),
        arrow_table_to_ipc_bytes(data),
    )

    print("created:", ".".join(table_id))
    print("location:", create_response.location)

    # テーブルを一覧表示して検証
    tables = ns.list_tables(ListTablesRequest(id=[DATABASE]))
    print("listed_after_create:", table_name in tables.tables)

    # データを読み取って検証
    desc = ns.describe_table(DescribeTableRequest(id=table_id))
    dataset = lance.dataset(desc.location, storage_options=desc.storage_options)
    result = dataset.to_table()

    print(result)

    # データ整合性チェック
    expected = data.to_pylist()
    actual = result.to_pylist()
    if actual != expected:
        raise AssertionError(f"readback mismatch: expected={expected}, actual={actual}")

    print("create_write_read: ok")


if __name__ == "__main__":
    main()

関連ドキュメント

データシリアル化ユーティリティ

lance_namespace モジュールの create_table インターフェイスには、IPC (プロセス間通信) バイトストリームとしてシリアル化された Arrow テーブルが必要です。

import pyarrow as pa


def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
    """PyArrow テーブルを IPC バイトストリームに変換"""
    sink = pa.BufferOutputStream()
    with pa.ipc.new_stream(sink, table.schema) as writer:
        writer.write_table(table)
    return sink.getvalue().to_pybytes()

DLF テーブルスキーマからのテストデータの構築

列名と型をハードコーディングする代わりに、DLF テーブルスキーマからフィールド情報を抽出してテストデータを構築します。

import pyarrow as pa
from lance_dlf.common.identifier import Identifier


def sample_value(field_type: str, row: int):
    """フィールドタイプごとにサンプル値を生成"""
    normalized = field_type.lower()
    if "int" in normalized:
        return row + 1
    if "string" in normalized or "char" in normalized or "varchar" in normalized:
        return f"value-{row + 1}"
    raise ValueError(f"Unsupported sample field type: {field_type}")


def build_sample_table(ns, database: str, table: str) -> pa.Table:
    """DLF テーブルスキーマからサンプルデータを構築"""
    raw_table = ns._api.get_table(Identifier(database, table))
    schema = raw_table.get_schema()
    fields = schema.fields if schema and schema.fields else []

    data = {}
    for field in fields:
        field_type = str(field.type)
        data[field.name] = [sample_value(field_type, row) for row in range(3)]

    return pa.table(data)

使用例:

desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
data = build_sample_table(ns, DATABASE, TABLE)

lance.write_dataset(
    data,
    desc.location,
    mode="overwrite",
    storage_options=desc.storage_options,
)