Panduan ini menjelaskan cara membaca dan menulis tabel Lance yang dikelola oleh DLF menggunakan mesin DataFrame Daft. Daft menyediakan API DataFrame lazy yang sangat cocok untuk skenario filtering kueri dan komputasi batch.
Untuk membaca dan menulis tabel Lance secara langsung dengan PyLance, lihat Bekerja dengan tabel Lance DLF menggunakan Python.
Istilah
Komponen | Peran |
DLF | Layanan katalog yang mengelola metadata database/tabel, menyimpan path tabel Lance, dan menerbitkan kredensial OSS temporary |
Lance/PyLance | Format data dan implementasi baca/tulis tingkat rendah yang bertanggung jawab atas I/O aktual pada set data Lance yang di-host di OSS |
Daft | Mesin komputasi DataFrame yang menyediakan antarmuka |
| Konektor yang mengambil path tabel dan kredensial OSS temporary dari DLF untuk digunakan oleh Daft. Hanya mengekspos tabel dengan |
Arsitektur
Kode pengguna
→ lance_namespace.connect("dlf", CONFIG) # Menghubungkan ke katalog DLF
→ DLF mengembalikan path tabel Lance + kredensial OSS temporary
→ apply_oss_environment(...) # Menyetel variabel lingkungan OSS_*
→ daft.read_lance("oss://...") # Membaca data
→ df.write_lance("oss://...", mode="append") # Menulis data Pemetaan konsep:
Database DLF → namespace Lance
Tabel DLF → tabel Lance
Prasyarat
Instal dependensi
python3 -m pip install lance-dlf daft lance-dlf secara otomatis menginstal lance_namespace, pyarrow, dan dependensi lain yang diperlukan.
Konfigurasi koneksi katalog
CONFIG = {
"uri": "http://<dlf-endpoint>", # Untuk akses publik, gunakan protokol 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>",
}Parameter konfigurasi:
Parameter | Deskripsi |
| Endpoint REST DLF Paimon. Gunakan protokol HTTPS untuk akses publik |
| Nama katalog DLF |
| Gunakan |
| ID wilayah DLF, misalnya |
| ID AccessKey untuk akses DLF |
| Rahasia AccessKey untuk akses DLF |
| (Opsional) Token keamanan untuk skenario STS |
| (Opsional) Endpoint publik OSS, misalnya |
ID dan rahasia AccessKey Anda merupakan kredensial penting untuk mengakses sumber daya Alibaba Cloud. Simpan dengan aman dan jangan pernah menyimpan AccessKey asli di repositori Git. Baca kredensial dari:
Variabel lingkungan
Sistem manajemen kunci
Konfigurasi runtime
Menghubungkan ke DLF
Mengimpor lance_dlf secara otomatis mendaftarkan namespace dlf:
import lance_namespace
import lance_dlf # noqa: F401
ns = lance_namespace.connect("dlf", CONFIG)
print(ns.namespace_id()) Membaca tabel yang sudah ada
Langkah 1: Dapatkan path tabel dan kredensial
Sebelum membaca atau menulis tabel dengan Daft, panggil describe_table() untuk mengambil path tabel dan kredensial OSS temporary dari DLF.
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 berisi dua bidang utama:
desc.location— path penyimpanan tabel Lance dalam formatoss://bucket/path/to/tabledesc.storage_options— kamus berisi kredensial OSS temporary
Langkah 2: Setel variabel lingkungan kredensial OSS
Daft menggunakan PyLance di balik layar untuk mengakses OSS. Definisikan helper apply_oss_environment yang memetakan kredensial yang dikeluarkan DLF ke variabel lingkungan OSS_*, lalu panggil fungsi tersebut:
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"]
# Terapkan kredensial
apply_oss_environment(desc.storage_options or {}) Langkah 3: Baca data tabel dengan Daft
import daft
df = daft.read_lance(desc.location)
df.show() Menulis data
Tambahkan ke tabel yang sudah ada
Setelah menyelesaikan Langkah 1: Dapatkan path tabel dan kredensial dan Langkah 2: Setel variabel lingkungan kredensial OSS di atas, tambahkan data dengan mode="append":
# Prasyarat: terhubung ke Katalog, dapatkan path/kredensial tabel, setel variabel lingkungan OSS
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
apply_oss_environment(desc.storage_options or {})
# Tambahkan data
append_df = daft.from_pydict({
"f0": [204],
"f1": ["daft-d"],
})
append_df.write_lance(desc.location, mode="append")
# Verifikasi penulisan
df2 = daft.read_lance(desc.location)
df2.show() Buat tabel baru dan tulis data
Tabel baru harus dibuat terlebih dahulu dengan ns.create_table() (yang juga menulis batch data awal). Gunakan Daft untuk pembacaan dan penulisan selanjutnya.
# Prasyarat: terhubung ke Katalog, dapatkan path/kredensial tabel, setel variabel lingkungan OSS
from datetime import datetime
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
# Serialisasi tabel Arrow ke byte 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()
# Buat tabel dengan data awal
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) Setelah pembuatan, ambil kredensial dengan describe_table dan gunakan Daft untuk membaca dan menulis:
# Dapatkan kredensial dan setel variabel lingkungan
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# Baca dan verifikasi
df = daft.read_lance(desc.location)
df.show()
# Tambahkan data dengan 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()
# Baca lagi untuk konfirmasi
appended_df = daft.read_lance(desc.location)
appended_df.show()
Output yang diharapkan:
[
{"f0": 201, "f1": "daft-a"},
{"f0": 202, "f1": "daft-b"},
{"f0": 203, "f1": "daft-c"},
{"f0": 204, "f1": "daft-d"}
]
Contoh lengkap
Skrip berikut menunjukkan alur kerja end-to-end: buat tabel baru → baca → tambahkan → verifikasi.
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>", # Hanya diperlukan untuk akses jaringan publik ke DLF
}
DATABASE = "default"
# Serialisasi tabel Arrow ke byte 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()
# Setel variabel lingkungan kredensial 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. Buat tabel baru
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. Dapatkan kredensial
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# 3. Baca dan verifikasi
read_df = daft.read_lance(desc.location)
read_df.show()
if df_to_pydict(read_df) != rows:
raise AssertionError("Initial readback mismatch")
# 4. Tambahkan data
append_rows = {
"f0": [204],
"f1": ["daft-d"],
}
append_df = daft.from_pydict(append_rows)
append_df.write_lance(desc.location, mode="append").show()
# 5. Verifikasi akhir
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()
Catatan penting
Inisialisasi tabel baru: Gunakan
ns.create_table(...)untuk membuat tabel baru dan menulis batch data pertama. Gunakan Daft untuk semua pembacaan dan penulisan selanjutnya.Pembersihan log: Kamus lengkap
storage_optionsberisi nilai AK/SK/token temporary. Cetak hanya daftar kuncinya demi keamanan:print(sorted((desc.storage_options or {}).keys()))