Este guia explica como ler e gravar tabelas Lance gerenciadas pelo DLF usando o mecanismo de DataFrame do Daft. O Daft fornece uma API de DataFrame lazy, ideal para cenários de filtragem de consultas e computação em lote.
Para ler e gravar tabelas Lance diretamente com o PyLance, consulte Work with DLF Lance using Python.
Terminologia
|
Componente |
Função |
|
DLF |
service de catálogo que gerencia metadados de bancos de dados e tabelas, armazena caminhos de tabelas Lance e emite credenciais temporárias do OSS |
|
Lance/PyLance |
Formato de dados e implementação de leitura e gravação de baixo nível, responsável pela E/S real em conjuntos de dados Lance hospedados no OSS |
|
Daft |
Mecanismo de computação de DataFrame que fornece as interfaces |
|
|
Conector que recupera caminhos de tabelas e credenciais temporárias do OSS a partir do DLF para uso pelo Daft. Expõe apenas tabelas com |
Arquitetura
User code
→ lance_namespace.connect("dlf", CONFIG) # Connect to DLF catalog
→ DLF returns Lance table path + temporary OSS credentials
→ apply_oss_environment(...) # Set OSS_* environment variables
→ daft.read_lance("oss://...") # Read data
→ df.write_lance("oss://...", mode="append") # Write data
Mapeamento de conceitos:
Banco de dados do DLF → namespace do Lance
Tabela do DLF → Tabela do Lance
Pré-requisitos
Instale as dependências
python3 -m pip install lance-dlf daft
O lance-dlf instala automaticamente o lance_namespace, o pyarrow e outras dependências necessárias.
Configuração mínima
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>",
}
Conecte-se ao DLF
A importação do lance_dlf registra automaticamente o namespace dlf:
import lance_namespace
import lance_dlf # noqa: F401
ns = lance_namespace.connect("dlf", CONFIG)
print(ns.namespace_id())
Leia uma tabela existente
Passo 1: Obtenha o caminho da tabela e as credenciais
Antes de ler ou gravar uma tabela com o Daft, chame describe_table() para recuperar o caminho da tabela e as credenciais temporárias do OSS a partir do 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()))
O objeto desc contém dois campos principais:
desc.location— o caminho de armazenamento da tabela Lance no formatooss://bucket/path/to/tabledesc.storage_options— um dicionário com as credenciais temporárias do OSS
Passo 2: Defina as variáveis de ambiente das credenciais do OSS
O Daft usa o PyLance internamente para acessar o OSS. Defina uma função auxiliar apply_oss_environment que mapeie as credenciais emitidas pelo DLF para as variáveis de ambiente OSS_* e, em seguida, chame-a:
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 credentials
apply_oss_environment(desc.storage_options or {})
Passo 3: Leia os dados da tabela com o Daft
import daft
df = daft.read_lance(desc.location)
df.show()
Grave dados
Acrescente dados a uma tabela existente
Após concluir o Step 1: Get the table path and credentials e o Step 2: Set OSS credential environment variables acima, acrescente dados usando mode="append":
# Prerequisites: connect to Catalog, get table path/credentials, set OSS env vars
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
apply_oss_environment(desc.storage_options or {})
# Append data
append_df = daft.from_pydict({
"f0": [204],
"f1": ["daft-d"],
})
append_df.write_lance(desc.location, mode="append")
# Verify the write
df2 = daft.read_lance(desc.location)
df2.show()
Crie uma nova tabela e grave dados
Novas tabelas devem ser criadas primeiro com ns.create_table() (que também grava o lote inicial de dados). Use o Daft para leituras e gravações subsequentes.
# Prerequisites: connect to Catalog, get table path/credentials, set OSS env vars
from datetime import datetime
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
# Serialize an Arrow table to IPC bytes
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()
# Create the table with initial data
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)
Após a criação, recupere as credenciais com describe_table e use o Daft para ler e gravar:
# Get credentials and set environment variables
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# Read and verify
df = daft.read_lance(desc.location)
df.show()
# Append data with 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()
# Read again to confirm
appended_df = daft.read_lance(desc.location)
appended_df.show()
Saída esperada:
[
{"f0": 201, "f1": "daft-a"},
{"f0": 202, "f1": "daft-b"},
{"f0": 203, "f1": "daft-c"},
{"f0": 204, "f1": "daft-d"}
]
Exemplo completo
O script a seguir demonstra o fluxo de trabalho de ponta a ponta: crie uma nova tabela → leia → acrescente → verifique.
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>", # Required only for public network access to DLF
}
DATABASE = "default"
# Serialize an Arrow table to IPC bytes
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()
# Set OSS credential environment variables
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. Create a new table
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. Get credentials
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# 3. Read and verify
read_df = daft.read_lance(desc.location)
read_df.show()
if df_to_pydict(read_df) != rows:
raise AssertionError("Initial readback mismatch")
# 4. Append 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. Final verification
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()
Notas importantes
Inicialização de nova tabela: Use
ns.create_table(...)para criar uma nova tabela e gravar o primeiro lote de dados. Utilize o Daft para todas as leituras e gravações subsequentes.Limpeza de logs: O dicionário completo
storage_optionscontém valores temporários de AK/SK/token. Por segurança, imprima apenas a lista de chaves:print(sorted((desc.storage_options or {}).keys()))