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
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
Python 3.10 ou superior.
-
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]" -
Instale o Daft.
pip install "daft>=0.7.17"
Configure os parâmetros
Prepare as seguintes informações:
|
Parâmetro |
Descrição |
|
|
Seu AccessKey ID. |
|
|
Seu AccessKey secret. |
|
|
ID da região onde o DLF está implantado, por exemplo, |
|
|
Nome do catalog no DLF (corresponde ao warehouse do Iceberg). |
|
|
Nome do banco de dados de destino (namespace do Iceberg). |
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.
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.