Este guia demonstra como usar o lance-dlf para conectar-se ao DLF Catalog, criar tabelas Lance, gravar dados e verificar resultados.
Funcionalidades
O módulo lance-dlf oferece os seguintes recursos principais:
Conecta-se ao DLF Catalog
Mapeia bancos de dados DLF para namespaces Lance (camada de mapeamento lógico)
Expõe apenas tabelas DLF do tipo
type=lance-tableObtém credenciais de acesso temporário ao OSS pela API
load_table_tokendo DLFConverte credenciais temporárias do OSS em
storage_optionspara o PyLance
O PyLance executa as operações reais de leitura e gravação de dados:
# Write data
lance.write_dataset(table, location, storage_options=storage_options)
# Read data
lance.dataset(location, storage_options=storage_options)
Instale lance-dlf
Instale o lance-dlf pelo PyPI:
python3 -m pip install lance-dlf
Configure conexão com o catálogo
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>",
}
Parâmetros de configuração
|
Parâmetro |
Descrição |
|
|
Endpoint do DLF REST Catalog, por exemplo, |
|
|
Nome do warehouse do DLF Catalog |
|
|
Use |
|
|
Região do DLF, como |
|
|
AccessKey ID para acesso ao DLF |
|
|
AccessKey Secret para acesso ao DLF |
|
|
(Opcional) Token de segurança para cenários STS |
|
|
(Opcional) Endpoint personalizado do OSS, como |
Importante: O AccessKey ID e o AccessKey Secret são credenciais críticas para acessar recursos da Alibaba Cloud. Armazene-as com segurança e nunca faça commit de AccessKeys reais em repositórios git. Leia as credenciais de:
Variáveis de ambiente
Um sistema de gerenciamento de chaves
Configurações de tempo de execução
Conectar-se ao DLF Catalog
Importar o módulo lance_dlf registra automaticamente a implementação do namespace dlf:
import lance_namespace
import lance_dlf # noqa: F401
# Connect to DLF Catalog
ns = lance_namespace.connect("dlf", CONFIG)
# Verify the connection
print(ns.namespace_id())
O método namespace_id() retorna informações sobre o catálogo conectado, incluindo o endpoint do DLF e o warehouse.
Visualize namespaces e tabelas
Mapeamento de conceitos
Bancos de dados DLF correspondem a namespaces Lance.
Exemplos de consulta
from lance_namespace import (
DescribeNamespaceRequest,
DescribeTableRequest,
ListNamespacesRequest,
ListTablesRequest,
)
DATABASE = "<database>"
TABLE = "<table>"
# List all namespaces
namespaces = ns.list_namespaces(ListNamespacesRequest(id=[]))
print(namespaces)
# Get namespace details
namespace = ns.describe_namespace(DescribeNamespaceRequest(id=[DATABASE]))
print(namespace)
# List all tables in the database
tables = ns.list_tables(ListTablesRequest(id=[DATABASE]))
print(tables)
# Get table details
table = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
print(table.location)
print(table.properties)
print(table.storage_options)
Campos retornados por describe_table
|
Campo |
Descrição |
|
|
Caminho físico de armazenamento do dataset Lance (geralmente |
|
|
Opções de esquema da tabela DLF; o campo |
|
|
Credenciais de acesso temporário para o PyLance ler e gravar no OSS |
Importante: O objeto storage_options contém AccessKey, AccessKey Secret e Token temporários. Sempre oculte esses valores nos logs.
Crie tabela Lance e gravar dados
Crie nova tabela
Se a tabela não existir, crie-a com a API de baixo nível:
Crie um esquema DLF a partir do esquema Arrow e gere os metadados da tabela no DLF.
Obtenha o local de armazenamento da tabela Lance (
location) após a criação.Gere credenciais de acesso temporário ao OSS com o
lance-dlf.Grave os dados Arrow no local especificado com o PyLance e notifique o DLF sobre o commit do Lance.
import lance
import pyarrow as pa
from lance_dlf.api.rest_exception import AlreadyExistsException
from lance_dlf.common.identifier import Identifier
from lance_dlf.schema.schema import Schema as DlfSchema
DATABASE = "default"
TABLE = "test_lance_create_001"
table_id = [DATABASE, TABLE]
# Build test data
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 table and write data
expected = DlfSchema.from_pyarrow_schema(data.schema, options={"type": "lance-table"})
identifier = Identifier(*table_id)
try:
ns._api.create_table(identifier, expected)
except AlreadyExistsException:
response = ns._api.get_table(identifier)
if response.get_schema() != expected:
raise ValueError("The existing table schema does not match the schema of the data to write")
else:
response = ns._api.get_table(identifier)
location = response.get_path()
storage_options = ns._build_storage_options(identifier)
lance.write_dataset(data, location, mode="overwrite", storage_options=storage_options)
ns.notify_lance_commit(table_id)
# Read and verify
dataset = lance.dataset(location, storage_options=storage_options)
result = dataset.to_table()
print(result)
Saída esperada
pyarrow.Table
f0: int64
f1: string
----
f0: [[101,102,103]]
f1: [["create-a","create-b","create-c"]]
Gravar em tabela vazia existente
Se já existir uma tabela vazia do tipo type=lance-table no DLF, obtenha o local de armazenamento e as credenciais de acesso por meio de describe_table e grave os dados com o PyLance.
import lance
import pyarrow as pa
from lance_namespace import DescribeTableRequest
DATABASE = "default"
TABLE = "test_lance_table"
table_id = [DATABASE, TABLE]
# Get table details
desc = ns.describe_table(DescribeTableRequest(id=table_id))
# Build test data
data = pa.table({
"f0": pa.array([1, 2, 3], type=pa.int64()),
"f1": pa.array(["value-1", "value-2", "value-3"], type=pa.string()),
})
# Write data
lance.write_dataset(
data,
desc.location,
mode="overwrite",
storage_options=desc.storage_options,
)
ns.notify_lance_commit(table_id)
# Read and verify
dataset = lance.dataset(desc.location, storage_options=desc.storage_options)
print(dataset.to_table())
Modos de gravação
|
Modo |
Descrição |
|
|
Substitui os dados existentes; útil para inicializar tabelas vazias ou de teste |
|
|
Adiciona dados; o esquema deve ser compatível |
Aviso: O modo overwrite substitui todos os dados existentes no dataset Lance. Tenha cuidado ao usar este modo.
Gerar dados de teste a partir do esquema da tabela DLF
Em vez de codificar nomes e tipos de colunas manualmente, extraia as informações de campo do esquema da tabela DLF para gerar dados de teste.
import pyarrow as pa
from lance_dlf.common.identifier import Identifier
def sample_value(field_type: str, row: int):
"""Generate sample value by field type"""
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:
"""Build sample data from DLF table schema"""
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)
Exemplo de uso
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,
)
Exemplo completo
Este exemplo demonstra o fluxo completo: conexão com o DLF, criação de tabela, gravação de dados, listagem de tabelas e verificação de resultados.
from datetime import datetime
import lance
import lance_namespace
import pyarrow as pa
from lance_dlf.api.rest_exception import AlreadyExistsException
from lance_dlf.common.identifier import Identifier
from lance_dlf.schema.schema import Schema as DlfSchema
from lance_namespace import 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 main():
# Connect to DLF Catalog
ns = lance_namespace.connect("dlf", CONFIG)
# Generate unique table name
table_name = "test_lance_create_" + datetime.now().strftime("%Y%m%d_%H%M%S")
table_id = [DATABASE, table_name]
# Build test data
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 table and write data
expected = DlfSchema.from_pyarrow_schema(data.schema, options={"type": "lance-table"})
identifier = Identifier(*table_id)
try:
ns._api.create_table(identifier, expected)
except AlreadyExistsException:
response = ns._api.get_table(identifier)
if response.get_schema() != expected:
raise ValueError("The existing table schema does not match the schema of the data to write")
else:
response = ns._api.get_table(identifier)
location = response.get_path()
storage_options = ns._build_storage_options(identifier)
lance.write_dataset(data, location, mode="overwrite", storage_options=storage_options)
ns.notify_lance_commit(table_id)
print("created:", ".".join(table_id))
print("location:", location)
# List tables and verify
tables = ns.list_tables(ListTablesRequest(id=[DATABASE]))
print("listed_after_create:", table_name in tables.tables)
# Read data and verify
dataset = lance.dataset(location, storage_options=storage_options)
result = dataset.to_table()
print(result)
# Data integrity check
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()
Execute o exemplo
python3 main.py