Todos os produtos
Search
Central de documentação

Data Lake Formation:Uso do DLF Lance com Python

Última atualização: Sep 17, 2026

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-table

  • Obtém credenciais de acesso temporário ao OSS pela API load_table_token do DLF

  • Converte credenciais temporárias do OSS em storage_options para 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

uri

Endpoint do DLF REST Catalog, por exemplo, http://...

warehouse

Nome do warehouse do DLF Catalog

token.provider

Use dlf para autenticação por AccessKey do DLF

dlf.region

Região do DLF, como cn-hangzhou

dlf.access-key-id

AccessKey ID para acesso ao DLF

dlf.access-key-secret

AccessKey Secret para acesso ao DLF

dlf.security-token

(Opcional) Token de segurança para cenários STS

dlf.oss-endpoint

(Opcional) Endpoint personalizado do OSS, como oss-cn-hangzhou.aliyuncs.com

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

location

Caminho físico de armazenamento do dataset Lance (geralmente oss://bucket/path)

properties

Opções de esquema da tabela DLF; o campo type deve ser lance-table

storage_options

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:

  1. Crie um esquema DLF a partir do esquema Arrow e gere os metadados da tabela no DLF.

  2. Obtenha o local de armazenamento da tabela Lance (location) após a criação.

  3. Gere credenciais de acesso temporário ao OSS com o lance-dlf.

  4. 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

overwrite

Substitui os dados existentes; útil para inicializar tabelas vazias ou de teste

append

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