Todos os produtos
Search
Central de documentação

Data Lake Formation:Use Daft com tabelas Iceberg do DLF

Última atualização: Sep 18, 2026

Daft é um motor de DataFrame distribuído de alto desempenho. Use o Daft para ler e gravar tabelas Apache Iceberg no Alibaba Cloud Data Lake Formation (DLF) e executar operações de DataFrame como filtragem e agregação.

Pré-requisitos

Nota

O serviço REST Iceberg do DLF é acessível somente de dentro de uma VPC. Execute o código deste tópico em um ambiente VPC na mesma região que o DLF (como uma instância ECS ou cluster EMR). Para o endpoint de cada região, consulte Iceberg REST endpoints.

Instale as dependências

  1. Python 3.10 ou superior.

  2. Instale o pacote PyIceberg compatível com DLF (pyiceberg-dlf).

    python3 -m venv venv
    source venv/bin/activate
    pip install -U pip
    # Uninstall pyiceberg, which cannot coexist with pyiceberg-dlf
    pip uninstall -y pyiceberg
    # rest-sigv4 is required (installs boto3 for REST sigv4 signing)
    pip install "pyiceberg-dlf[rest-sigv4,pyarrow,pandas]"
  3. Instale o Daft.

    pip install "daft>=0.7.17"

Configure os parâmetros

Prepare as seguintes informações:

Parâmetro

Descrição

${accessKeyId}

Seu AccessKey ID.

${accessKeySecret}

Seu AccessKey secret.

${regionId}

ID da região onde o DLF está implantado, por exemplo, ap-southeast-1. Para a lista de IDs de região, consulte Regions and endpoints.

${catalogName}

Nome do catalog no DLF (corresponde ao warehouse do Iceberg).

${database}

Nome do banco de dados de destino (namespace do Iceberg).

Importante

Proteja suas credenciais AccessKey. Não as insira diretamente no código nem as inclua em repositórios de código. Recupere-as a partir de variáveis de ambiente ou de um serviço de gerenciamento de segredos.

Conexão e inicialização

Conecte-se ao catalog do DLF

Conecte-se ao catalog do DLF via protocolo REST do Iceberg usando o load_catalog do PyIceberg.

from pyiceberg.catalog import load_catalog

REGION = "${regionId}"

catalog = load_catalog(
    "dlf",
    **{
        "type": "rest",
        "uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
        "warehouse": "${catalogName}",
        "rest.signing-name": "DlfNext",
        "rest.signing-region": REGION,
        "rest.sigv4-enabled": "true",
        "client.access-key-id": "${accessKeyId}",
        "client.secret-access-key": "${accessKeySecret}",
        "client.region": REGION,
        "s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com",
    },
)

Crie ou carregue uma tabela

Crie uma nova tabela

Crie uma tabela Iceberg usando PyIceberg. O DLF gerencia os metadados da tabela.

from pyiceberg.schema import Schema
from pyiceberg.types import LongType, StringType, DoubleType, NestedField
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC

schema = Schema(
    NestedField(field_id=1, name="id", field_type=LongType(), required=True),
    NestedField(field_id=2, name="name", field_type=StringType(), required=False),
    NestedField(field_id=3, name="category", field_type=StringType(), required=False),
    NestedField(field_id=4, name="value", field_type=DoubleType(), required=False),
)

table = catalog.create_table(
    identifier=("${database}", "daft_demo_table"),
    schema=schema,
    partition_spec=UNPARTITIONED_PARTITION_SPEC,
)
print("Table created. Data location:", table.location())

Carregue uma tabela existente

Para trabalhar com uma tabela já existente, chame catalog.load_table().

table = catalog.load_table(("${database}", "table_name"))

Operações de dados

Grave dados

Construa um DataFrame com o Daft e grave-o em uma tabela Iceberg. O write_iceberg retorna uma tabela resumida dos arquivos de dados gravados nessa operação.

Nota

O Daft 0.7.15 configura automaticamente o acesso ao OSS para tabelas Iceberg em caminhos oss://; portanto, write_iceberg e read_iceberg não exigem configuração manual de IOConfig.

import daft

df = daft.from_pydict({
    "id": [1, 2, 3, 4, 5],
    "name": ["item_1", "item_2", "item_3", "item_4", "item_5"],
    "category": ["A", "B", "A", "B", "A"],
    "value": [10.5, 21.0, 31.5, 42.0, 52.5],
})

result = df.write_iceberg(table, mode="append")
result.show()

O parâmetro mode aceita "append" (acrescentar dados) e "overwrite" (sobrescrever dados existentes).

Leia dados

Após gravar os dados, recarregue a tabela para obter o snapshot mais recente e atualizar as credenciais STS; em seguida, leia-a com o Daft.

table = catalog.load_table(("${database}", "daft_demo_table"))

df = daft.read_iceberg(table)
df.show()

O daft.read_iceberg é lazy: retorna um handle de DataFrame e constrói um plano de execução. Os dados são lidos do OSS apenas quando você chama ações como show() ou collect().

Transforme dados

Com o DataFrame lido da tabela (contendo as cinco linhas de exemplo), é possível executar filtragem, criação de colunas derivadas, agregação e ordenação.

df = daft.read_iceberg(table).collect()

# Filter: rows where value is greater than 30
df.where(df["value"] > 30).show()

# Derived column: add value_x2 = value * 2
df.with_column("value_x2", df["value"] * 2).show()

# Group and aggregate: count rows and sum values by category
df.groupby("category").agg(
    daft.col("id").count().alias("row_count"),
    daft.col("value").sum().alias("value_sum"),
).sort("category").show()

# Sort: descending by value
df.sort("value", desc=True).show()

Exemplo completo

O código a seguir reúne todas as etapas anteriores e pode ser copiado e executado diretamente.

import uuid

import daft
from pyiceberg.catalog import load_catalog
from pyiceberg.schema import Schema
from pyiceberg.types import LongType, StringType, DoubleType, NestedField
from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC

# ==================== Configuration ====================
ACCESS_KEY_ID = "${accessKeyId}"
ACCESS_KEY_SECRET = "${accessKeySecret}"
REGION = "${regionId}"
CATALOG_NAME = "${catalogName}"
DATABASE = "${database}"
# =======================================================

def create_catalog():
    """Connect to the DLF Iceberg REST Catalog."""
    return load_catalog(
        "dlf",
        **{
            "type": "rest",
            "uri": f"http://{REGION}-vpc.dlf.aliyuncs.com/iceberg",
            "warehouse": CATALOG_NAME,
            "rest.signing-name": "DlfNext",
            "rest.signing-region": REGION,
            "rest.sigv4-enabled": "true",
            "client.access-key-id": ACCESS_KEY_ID,
            "client.secret-access-key": ACCESS_KEY_SECRET,
            "client.region": REGION,
            "s3.endpoint": f"https://oss-{REGION}-internal.aliyuncs.com",
        },
    )

def create_table(catalog, table_name):
    """Create an unpartitioned Iceberg table using PyIceberg."""
    schema = Schema(
        NestedField(field_id=1, name="id", field_type=LongType(), required=True),
        NestedField(field_id=2, name="name", field_type=StringType(), required=False),
        NestedField(field_id=3, name="category", field_type=StringType(), required=False),
        NestedField(field_id=4, name="value", field_type=DoubleType(), required=False),
    )
    return catalog.create_table(
        identifier=(DATABASE, table_name),
        schema=schema,
        partition_spec=UNPARTITIONED_PARTITION_SPEC,
    )

def main():
    catalog = create_catalog()
    table_name = f"daft_demo_{uuid.uuid4().hex[:8]}"
    table = None
    try:
        # 1. Create table (PyIceberg)
        table = create_table(catalog, table_name)
        print("Table created. Data location:", table.location())

        # 2. Write data using Daft
        df = daft.from_pydict({
            "id": [1, 2, 3, 4, 5],
            "name": ["item_1", "item_2", "item_3", "item_4", "item_5"],
            "category": ["A", "B", "A", "B", "A"],
            "value": [10.5, 21.0, 31.5, 42.0, 52.5],
        })
        df.write_iceberg(table, mode="append").show()

        # 3. Reload table to get latest snapshot, then read using Daft
        table = catalog.load_table((DATABASE, table_name))
        result = daft.read_iceberg(table).collect()
        result.show()

        # 4. DataFrame transformations
        result.groupby("category").agg(
            daft.col("id").count().alias("row_count"),
            daft.col("value").sum().alias("value_sum"),
        ).sort("category").show()
    finally:
        # 5. Drop table (PyIceberg)
        if table is not None:
            catalog.drop_table((DATABASE, table_name))
            print("Table dropped:", table_name)

if __name__ == "__main__":
    main()

A saída esperada da etapa de gravação (a tabela de resultado de write_iceberg) é semelhante ao seguinte:

╭───────────┬───────┬───────────┬────────────────────────────────╮
│ operation ┆ rows  ┆ file_size ┆ file_name                      │
╞═══════════╪═══════╪═══════════╪════════════════════════════════╡
│ ADD       ┆ 5     ┆ 1708      ┆ oss://<bucket>/.../xxx.parquet │
╰───────────┴───────┴───────────┴────────────────────────────────╯

Apêndice

Terminologia

Componente

Descrição

DLF

Data Lake Formation. Fornece serviços de Iceberg REST Catalog e gerenciamento unificado de metadados de tabelas.

PyIceberg

Cliente Python para Apache Iceberg. Conecta-se aos catalogs do DLF, cria e remove tabelas, e confirma transações.

Daft

Motor de DataFrame distribuído. Grava dados em tabelas Iceberg, lê dados dessas tabelas e executa transformações de DataFrame.

OSS

Object Storage Service. Camada de armazenamento físico dos arquivos de dados de tabelas Iceberg (formato Parquet).

Arquitetura

Daft e PyIceberg trabalham em conjunto com uma separação clara de responsabilidades:

  • Plano de controle (PyIceberg): Conecta-se ao DLF usando o protocolo REST do Iceberg. Gerencia conexões com o catalog, criação e remoção de tabelas, e commits de snapshot.

  • Plano de dados (Daft): Os arquivos de dados das tabelas Iceberg ficam armazenados no OSS. O Daft lê e grava esses arquivos Parquet em paralelo e oferece capacidades de computação sobre DataFrame (filtragem, colunas derivadas, agregação, ordenação, entre outras).

Nas gravações, o Daft escreve os arquivos de dados Parquet e confirma snapshots Iceberg de forma atômica por meio do PyIceberg. Nas leituras, o Daft aproveita os metadados do Iceberg para realizar pruning de partições e filtragem de arquivos.