このガイドでは、Daft DataFrame エンジンを使用して、DLF が管理する Lance テーブルを読み書きする方法について説明します。Daft は、クエリのフィルタリングやバッチ計算のシナリオに適した遅延評価の DataFrame API を提供します。
PyLance を使用して Lance テーブルを直接読み書きするには、「Python を使用した DLF Lance テーブルの操作」をご参照ください。
用語
コンポーネント | 役割 |
DLF | データベース/テーブルのメタデータを管理し、Lance テーブルのパスを格納し、一時的な OSS の認証情報を発行するカタログサービス |
Lance/PyLance | OSS 上でホストされている Lance データセットに対する実際の I/O を担当するデータフォーマットと低レベルの読み書き実装 |
Daft |
|
| Daft が DLF からテーブルパスと一時的な OSS 認証情報を取得するためのコネクタ。 |
アーキテクチャ
User code
→ lance_namespace.connect("dlf", CONFIG) # DLF カタログに接続します
→ DLF returns Lance table path + temporary OSS credentials # DLF は Lance テーブルのパスと一時的な OSS 認証情報を返します
→ apply_oss_environment(...) # OSS_* 環境変数を設定します
→ daft.read_lance("oss://...") # データを読み取ります
→ df.write_lance("oss://...", mode="append") # データを書き込みます 概念のマッピング:
DLF データベース → Lance 名前空間
DLF テーブル → Lance テーブル
前提条件
依存関係のインストール
python3 -m pip install lance-dlf daft lance-dlf は、lance_namespace、pyarrow、およびその他の必要な依存関係を自動的にインストールします。
カタログ接続の設定
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
ns = lance_namespace.connect("dlf", CONFIG)
print(ns.namespace_id()) 既存テーブルの読み取り
ステップ 1:テーブルパスと認証情報の取得
Daft でテーブルを読み書きする前に、describe_table() を呼び出して、DLF からテーブルパスと一時的な OSS 認証情報を取得します。
from lance_namespace import DescribeTableRequest
DATABASE = "<database>"
TABLE = "<table>"
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
print(desc.location)
print(sorted((desc.storage_options or {}).keys())) desc には、次の 2 つの主要なフィールドが含まれています:
desc.location—oss://bucket/path/to/table形式の Lance テーブルのストレージパスdesc.storage_options— 一時的な OSS 認証情報のディクショナリ
ステップ 2:OSS 認証情報環境変数の設定
Daft は内部で PyLance を使用して OSS にアクセスします。DLF が発行した認証情報を OSS_* 環境変数にマッピングするため、ヘルパー関数 apply_oss_environment を定義して呼び出します。
import os
def apply_oss_environment(storage_options: dict) -> None:
os.environ["OSS_ENDPOINT"] = storage_options["oss_endpoint"]
os.environ["OSS_ACCESS_KEY_ID"] = storage_options["oss_access_key_id"]
os.environ["OSS_ACCESS_KEY_SECRET"] = storage_options["oss_secret_access_key"]
if storage_options.get("oss_security_token"):
os.environ["OSS_SECURITY_TOKEN"] = storage_options["oss_security_token"]
if storage_options.get("oss_region"):
os.environ["OSS_REGION"] = storage_options["oss_region"]
# 認証情報を適用
apply_oss_environment(desc.storage_options or {}) ステップ 3:Daft でのテーブルデータの読み取り
import daft
df = daft.read_lance(desc.location)
df.show() データの書き込み
既存テーブルへの追加
上記の「ステップ 1:テーブルパスと認証情報の取得」と「ステップ 2:OSS 認証情報環境変数の設定」を完了した後、mode="append" を使用してデータを追加します:
# 前提条件:カタログへの接続、テーブルパス/認証情報の取得、OSS 環境変数の設定
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
apply_oss_environment(desc.storage_options or {})
# データの追加
append_df = daft.from_pydict({
"f0": [204],
"f1": ["daft-d"],
})
append_df.write_lance(desc.location, mode="append")
# 書き込みの検証
df2 = daft.read_lance(desc.location)
df2.show() 新規テーブルの作成とデータの書き込み
新しいテーブルを作成するには、まず ns.create_table() を使用する必要があります (これにより、データの初期バッチも書き込まれます)。その後の読み書きには Daft を使用します。
# 前提条件:カタログへの接続、テーブルパス/認証情報の取得、OSS 環境変数の設定
from datetime import datetime
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
# Arrow テーブルを IPC バイトにシリアル化
def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, table.schema) as writer:
writer.write_table(table)
return sink.getvalue().to_pybytes()
# 初期データでテーブルを作成
table_name = "test_lance_daft_" + datetime.now().strftime("%Y%m%d_%H%M%S")
table_id = [DATABASE, table_name]
rows = {
"f0": [201, 202, 203],
"f1": ["daft-a", "daft-b", "daft-c"],
}
arrow_table = pa.table(rows)
create_response = ns.create_table(
CreateTableRequest(id=table_id),
arrow_table_to_ipc_bytes(arrow_table),
)
print(create_response.location) 作成後、describe_table で認証情報を取得し、Daft を使用して読み書きします:
# 認証情報を取得して環境変数を設定
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# 読み取りと検証
df = daft.read_lance(desc.location)
df.show()
# Daft でデータを追加
append_rows = {
"f0": [204],
"f1": ["daft-d"],
}
append_df = daft.from_pydict(append_rows)
meta = append_df.write_lance(desc.location, mode="append")
meta.show()
# 再度読み取って確認
appended_df = daft.read_lance(desc.location)
appended_df.show()
期待される出力:
[
{"f0": 201, "f1": "daft-a"},
{"f0": 202, "f1": "daft-b"},
{"f0": 203, "f1": "daft-c"},
{"f0": 204, "f1": "daft-d"}
]
完全な例
次のスクリプトは、新規テーブルの作成 → 読み取り → 追加 → 検証という、エンドツーエンドのワークフローを示しています。
from __future__ import annotations
from datetime import datetime
import os
import daft
import lance_namespace
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
import lance_dlf # noqa: F401
CONFIG = {
"uri": "http://<DLF-ENDPOINT>",
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": "<ACCESS-KEY-ID>",
"dlf.access-key-secret": "<ACCESS-KEY-SECRET>",
"dlf.oss-endpoint": "<OSS-ENDPOINT>", # DLF へのパブリックネットワークアクセスにのみ必要
}
DATABASE = "default"
# Arrow テーブルを IPC バイトにシリアル化
def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, table.schema) as writer:
writer.write_table(table)
return sink.getvalue().to_pybytes()
# OSS 認証情報環境変数の設定
def apply_oss_environment(storage_options: dict) -> None:
os.environ["OSS_ENDPOINT"] = storage_options["oss_endpoint"]
os.environ["OSS_ACCESS_KEY_ID"] = storage_options["oss_access_key_id"]
os.environ["OSS_ACCESS_KEY_SECRET"] = storage_options["oss_secret_access_key"]
if storage_options.get("oss_security_token"):
os.environ["OSS_SECURITY_TOKEN"] = storage_options["oss_security_token"]
if storage_options.get("oss_region"):
os.environ["OSS_REGION"] = storage_options["oss_region"]
def df_to_pydict(df):
try:
return df.to_pydict()
except AttributeError:
return df.collect().to_pydict()
def main() -> None:
ns = lance_namespace.connect("dlf", CONFIG)
# 1. 新規テーブルの作成
table_name = "test_lance_daft_" + datetime.now().strftime("%Y%m%d_%H%M%S")
table_id = [DATABASE, table_name]
rows = {
"f0": [201, 202, 203],
"f1": ["daft-a", "daft-b", "daft-c"],
}
arrow_table = pa.table(rows)
create_response = ns.create_table(
CreateTableRequest(id=table_id),
arrow_table_to_ipc_bytes(arrow_table),
)
print("created:", ".".join(table_id))
print("location:", create_response.location)
# 2. 認証情報の取得
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# 3. 読み取りと検証
read_df = daft.read_lance(desc.location)
read_df.show()
if df_to_pydict(read_df) != rows:
raise AssertionError("Initial readback mismatch")
# 4. データの追加
append_rows = {
"f0": [204],
"f1": ["daft-d"],
}
append_df = daft.from_pydict(append_rows)
append_df.write_lance(desc.location, mode="append").show()
# 5. 最終検証
appended_df = daft.read_lance(desc.location)
appended_df.show()
expected = {
"f0": rows["f0"] + append_rows["f0"],
"f1": rows["f1"] + append_rows["f1"],
}
if df_to_pydict(appended_df) != expected:
raise AssertionError("Daft append readback mismatch")
print("daft + dlf + lance: ok")
if __name__ == "__main__":
main()
重要な注意点
新規テーブルの初期化:
ns.create_table(...)を使用して新規テーブルを作成し、データの最初のバッチを書き込みます。その後のすべての読み書きには Daft を使用します。ログのサニタイズ:完全な
storage_optionsディクショナリには、一時的な AccessKey ID、AccessKey シークレット、セキュリティトークンの値が含まれています。安全のために、キーのリストのみを出力します:print(sorted((desc.storage_options or {}).keys()))