Todos os produtos
Search
Central de documentação

Data Lake Formation:Trabalhe com tabelas Lance do DLF usando o Daft

Última atualização: Sep 18, 2026

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.

Nota

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 read_lance / write_lance

lance-dlf

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 type=lance-table

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 formato oss://bucket/path/to/table

  • desc.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_options contém valores temporários de AK/SK/token. Por segurança, imprima apenas a lista de chaves: print(sorted((desc.storage_options or {}).keys()))