このガイドでは、lance-dlf を使用して、DLF カタログに接続し、Lance テーブルを作成し、データを書き込み、結果を検証する方法について説明します。
Daft を使用して DLF の Lance テーブルを操作するには、「Daft を使用して DLF Lance テーブルを操作する」をご参照ください。
機能概要
lance-dlf は、以下のコア機能を提供します:
DLF カタログへの接続
DLF データベースを Lance 名前空間にマッピング (論理マッピングレイヤー)
type=lance-tableの DLF テーブルのみを公開しますDLF の
load_table_tokenAPI 経由で一時的な 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>",
}設定パラメーター:
パラメーター | 説明 |
| DLF Paimon REST エンドポイント。パブリックアクセスには HTTPS プロトコルを使用します |
| DLF カタログ名 |
|
|
| DLF リージョン ID、例: |
| DLF アクセス用の AccessKey ID |
| DLF アクセス用の AccessKey Secret |
| (オプション) STS シナリオ用のセキュリティトークン |
| (任意) OSS パブリック エンドポイント (例: |
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 テーブルの作成とデータの書き込み
テーブルが存在しない場合、テーブルを作成してデータを挿入します。メカニズム:
ns.create_table()を使用してテーブルを作成しますDLF から Lance テーブルのストレージの場所(
location)を取得するlance-dlf経由で一時的な OSS アクセス資格情報を取得する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())書き込みモード:
モード | 説明 |
| 既存のデータを上書きします。空のテーブルやテストテーブルの初期化に便利です。 |
| データを追記します。スキーマに互換性がある必要があります。 |
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 の戻り値フィールド:
フィールド | 説明 |
| Lance データセットの物理ストレージパス (通常は |
| DLF テーブルスキーマのオプションでは、 |
| 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,
)