Este artigo explica como usar o PyPaimon para criar, ler e gravar dados em tabelas Paimon no Data Lake Formation (DLF).
Integração entre PyPaimon e DLF
O PyPaimon é o SDK Python para o Apache Paimon. Ele oferece recursos eficientes de ingestão de dados, permitindo que você leia, grave e processe dados de tabelas Paimon diretamente com Python.
Integração com o catálogo DLF
Ao importar o pacote de extensão pypaimon_dlf2 e configurar um catálogo DLF, os metadados das tabelas Paimon são sincronizados automaticamente com o Alibaba Cloud Data Lake Formation (DLF).
Essa integração oferece os seguintes benefícios principais:
Interoperabilidade entre múltiplos engines: Após hospedar os metadados no DLF, outros engines de computação do Alibaba Cloud — como MaxCompute, Hologres e Alibaba Cloud EMR — podem acessar os dados Paimon de forma transparente.
Governança unificada: Os recursos de gerenciamento de data lake do DLF permitem controlar o ciclo de vida das tabelas Paimon e otimizar automaticamente seus formatos de armazenamento.
Pré-requisitos
Você possui um DLF data catalog.
A versão do Python deve ser 3,8 ou superior. Execute
python3 --versionpara verificar a versão atual.
Procedimento
Etapa 1: Prepare o ambiente
-
Execute o comando a seguir para instalar o SDK PyPaimon.
pip3 install pypaimon==1.4.1 (Opcional) Após a instalação, execute o comando
pip3 show pypaimonpara verificar se a instalação foi concluída com êxito.
Etapa 2: Acesse uma tabela Paimon no DLF
-
No diretório de destino, execute o comando a seguir para criar um novo arquivo chamado
testdlf.py.vim testdlf.py -
No arquivo
testdlf.py, adicione o código de exemplo completo a seguir. Este exemplo demonstra como criar uma tabela Paimon no DLF e, em seguida, ler e gravar dados nela. Para detalhes sobre configurações de parâmetros e outros métodos de leitura e gravação, consulte Code details.import pyarrow as pa import pandas as pd from pypaimon import CatalogFactory from pypaimon import Schema # Create a catalog. catalog_options = { 'metastore': 'rest', 'uri': "http://${region_id}-vpc.dlf.aliyuncs.com", 'warehouse': "${catalog_name}", 'dlf.region': '${region_id}', "token.provider": "dlf", 'dlf.access-key-id': "xxx", 'dlf.access-key-secret': "xxxx", } catalog = CatalogFactory.create(catalog_options) # Create a database. catalog.create_database( name='testdb', ignore_if_exists=True # Specifies whether to ignore the error if the database already exists. ) # Create a schema. pa_schema = pa.schema([ ('date', pa.string()), ('hour', pa.string()), ('key', pa.int64()), ('value', pa.string()) ]) schema = Schema.from_pyarrow_schema( pa_schema=pa_schema, partition_keys=['date', 'hour'], primary_keys=['date', 'hour', 'key'], options={'bucket': '2'}, comment='my test table' ) # Create a table. catalog.create_table( identifier='testdb.tb', schema=schema, ignore_if_exists=True # Specifies whether to ignore the error if the table already exists. ) table = catalog.get_table('testdb.tb') # Create table write and commit operations. write_builder = table.new_batch_write_builder() table_write = write_builder.new_write() table_commit = write_builder.new_commit() # Write data to the table. Both PyArrow and Pandas are supported. # Write sample Pandas data. data = { 'date': ['2024-12-01', '2024-12-01', '2024-12-02'], 'hour': ['08', '09', '08'], 'key': [1, 2, 3], 'value': ['AAA', 'BBB', 'CCC'], } dataframe = pd.DataFrame(data) table_write.write_pandas(dataframe) # Commit the data. table_commit.commit(table_write.prepare_commit()) # Close the resources. table_write.close() table_commit.close() # Read data from the table. Multiple output formats are supported. read_builder = table.new_read_builder() predicate_builder = read_builder.new_predicate_builder() predicate = predicate_builder.equal('date', '2024-12-01') read_builder = read_builder.with_filter(predicate) table_scan = read_builder.new_scan() splits = table_scan.plan().splits() table_read = read_builder.new_read() pa_table = table_read.to_arrow(splits) print(pa_table)
Etapa 3: Execute o arquivo Python
Acesse o diretório de destino e execute o comando a seguir para rodar o script Python.
python3 testdlf.py
O retorno esperado é o seguinte.
root@iZxxx:/opt# python3 testdlf.py
pyarrow.Table
date: string not null
hour: string not null
key: int64 not null
value: string
----
date: [["2024-12-01"],["2024-12-01"]]
hour: [["09"],["08"]]
key: [[2],[1]]
value: [["BBB"],["AAA"]]
Detalhes do código
Criar uma tabela Paimon no DLF
-
Crie um catálogo DLF Paimon.
NotaÉ necessário criar um catálogo para acessar tabelas Paimon no DLF.
# The catalog_options is a dictionary where both keys and values are strings. catalog_options = { 'metastore': 'rest', 'uri': "http://${region_id}-vpc.dlf.aliyuncs.com", 'warehouse': "${catalog_name}", 'dlf.region': '${region_id}', "token.provider": "dlf", 'dlf.access-key-id': "xxx", 'dlf.access-key-secret': "xxxx", } catalog = CatalogFactory.create(catalog_options)A tabela a seguir descreve os parâmetros.
Parâmetro
Descrição
metastore
Defina como o valor fixo
rest, que indica a conexão com o DLF por meio do protocolo REST Catalog.dlf.region
O ID da região do DLF. Para mais informações, consulte Endpoints.
uri
O endpoint do DLF REST Catalog. Em ambiente VPC, use
http://${region_id}-vpc.dlf.aliyuncs.com. Em rede pública, usehttps://dlfnext.${region_id}.aliyuncs.com. Para mais informações, consulte Endpoints.warehouse
O nome do catálogo de dados do DLF. Você pode visualizar o nome no console do Data Lake Formation. Para mais informações, consulte Data catalog.
dlf.access-key-id
O AccessKey ID necessário para acessar o service DLF. Para mais informações, consulte Create an AccessKey pair.
dlf.access-key-secret
O AccessKey secret necessário para acessar o service DLF. Para mais informações, consulte Create an AccessKey pair.
token.provider
Defina como o valor fixo
dlf, que indica que o service DLF fornece o token de acesso.max-workers
Opcional. O número de threads simultâneas para leitura de dados no PyPaimon. Deve ser um inteiro maior ou igual a 1. O padrão é 1, que indica leitura serial.
-
Crie um banco de dados.
Em um catálogo Paimon, toda tabela pertence a um banco de dados específico. Crie bancos de dados para organizar e gerenciar suas tabelas.
catalog.create_database( name='database_name', ignore_if_exists=True, # Specifies whether to ignore the error if the database already exists. properties={'key': 'value'} # Optional. The database properties. ) -
Crie um schema.
Um schema inclui definições de colunas, chaves de partição, chaves primárias, opções de tabela e comentários. As definições de colunas são descritas usando
pyarrow.Schema. Os demais parâmetros são opcionais. Você pode construir umpyarrow.Schemade uma das seguintes formas.PyArrow
Use o método
pyarrow.schema. O código a seguir apresenta um exemplo.import pyarrow as pa from pypaimon import Schema pa_schema = pa.schema([ ('date', pa.string()), ('hour', pa.string()), ('key', pa.int64()), ('value', pa.string()) ]) schema = Schema( pa_schema=pa_schema, partition_keys=['date', 'hour'], primary_keys=['date', 'hour', 'key'], options={'bucket': '2'}, comment='my test table' )NotaPara informações sobre o mapeamento de tipos de dados entre
pyarrowePaimon, consulte PyPaimon data type mapping.Pandas
Se você já possui dados em Pandas, derive o schema diretamente de um
pandas.DataFrame. O código a seguir apresenta um exemplo.import pandas as pd import pyarrow as pa from pypaimon import Schema # This is sample DataFrame data. data = { 'date': ['2024-12-01', '2024-12-01', '2024-12-02'], 'hour': ['08', '09', '08'], 'key': [1, 2, 3], 'value': ['AAA', 'BBB', 'CCC'], } dataframe = pd.DataFrame(data) # Obtain the pyarrow.Schema from the DataFrame. record_batch = pa.RecordBatch.from_pandas(dataframe) pa_schema = record_batch.schema schema = Schema( pa_schema=pa_schema, partition_keys=['date', 'hour'], primary_keys=['date', 'hour', 'key'], options={'bucket': '2'}, comment='my test table' ) -
Crie e obtenha uma tabela.
catalog.create_table( identifier='database_name.table_name', schema=schema, ignore_if_exists=True # Specifies whether to ignore the error if the table already exists. ) table = catalog.get_table('database_name.table_name')
Gravar dados em uma tabela
O PyPaimon não oferece suporte, no momento, à gravação de dados em tabelas de chave primária com a opção bucket definida como -1.
-
Crie as operações de gravação e commit da tabela.
# Create table write and commit operations. write_builder = table.new_batch_write_builder() table_write = write_builder.new_write() table_commit = write_builder.new_commit() # Write sample Pandas data. data = { 'date': ['2024-12-01', '2024-12-01', '2024-12-02'], 'hour': ['08', '09', '08'], 'key': [1, 2, 3], 'value': ['AAA', 'BBB', 'CCC'], } -
Grave os dados na tabela por um dos seguintes métodos:
Para grandes volumes de dados, use PyArrow. Para conjuntos menores — normalmente alguns gigabytes ou menos — o Pandas pode ser mais eficiente.
PyArrow
Grave os dados como um
pyarrow.Tableou umpyarrow.RecordBatch. O pyarrow.RecordBatch é mais adequado para processamento em stream.-
Método 1: Grave um pyarrow.Table
# Create fields. fields = [ pa.field('date', pa.string()), pa.field('hour', pa.string()), pa.field('key', pa.int64()), pa.field('value', pa.string()) ] # Create a schema from the fields. schema = pa.schema(fields) # Create a table. pa_table = pa.Table.from_arrays(data, schema) # Write the data. table_write.write_arrow(pa_table) -
Método 2: Grave um pyarrow.RecordBatch
# Create fields. fields = [ pa.field('date', pa.string()), pa.field('hour', pa.string()), pa.field('key', pa.int64()), pa.field('value', pa.string()) ] # Create a schema from the fields. schema = pa.schema(fields) # Create a RecordBatch. record_batch = pa.RecordBatch.from_arrays(data, schema) # Write the data. table_write.write_arrow_batch(record_batch)
Pandas
Grave os dados a partir de um pandas.DataFrame.
import pandas as pd dataframe = pd.DataFrame(data) table_write.write_pandas(dataframe) -
-
Faça o commit dos dados e libere os recursos.
# Commit the data. table_commit.commit(table_write.prepare_commit()) # Close the resources. table_write.close() table_commit.close()
Ler dados de uma tabela
-
Crie um ReadBuilder.
read_builder = table.new_read_builder() -
Use um PredicateBuilder para construir e aplicar condições de filtro.
-
Por exemplo, leia apenas os dados em que
dateseja2024-12-01.predicate_builder = read_builder.new_predicate_builder() predicate = predicate_builder.equal('date', '2024-12-01') read_builder = read_builder.with_filter(predicate) -
Por exemplo, projete apenas as colunas
keyevalue.read_builder = read_builder.with_projection(['key', 'value'])
NotaPara mais informações sobre as condições de filtro suportadas, consulte PyPaimon filter conditions.
-
-
Obtenha os
splits.table_scan = read_builder.new_scan() splits = table_scan.plan().splits() -
Converta os
splitspara diferentes formatos de saída.Apache Arrow
-
Leia todos os dados em um
pyarrow.Table.table_read = read_builder.new_read() pa_table = table_read.to_arrow(splits) print(pa_table) # Sample output: # pyarrow.Table # key: int64 not null # value: string # ---- # key: [[2],[1]] # value: [["BBB"],["AAA"]] -
Leia os dados em um
pyarrow.RecordBatchReadere itere sobre os batches.table_read = read_builder.new_read() for batch in table_read.to_arrow_batch_reader(splits): print(batch) # Sample output: # pyarrow.RecordBatch # key: int64 # value: string # ---- # key: [1,2] # value: ["AAA","BBB"]
Pandas
Leia os dados em um
pandas.DataFrame.table_read = read_builder.new_read() df = table_read.to_pandas(splits) print(df) # Sample output: # key value # 0 1 AAA # 1 2 BBBDuckDB
ImportanteO DuckDB precisa estar instalado. Execute
pip install duckdbpara instalá-lo.Converta os dados para uma tabela DuckDB em memória e faça consultas sobre ela.
table_read = read_builder.new_read() duckdb_con = table_read.to_duckdb(splits, 'duckdb_table') print(duckdb_con.query("SELECT * FROM duckdb_table").fetchdf()) # Sample output: # key value # 0 1 AAA # 1 2 BBB print(duckdb_con.query("SELECT * FROM duckdb_table WHERE key = 1").fetchdf()) # Sample output: # key value # 0 1 AAARay
ImportanteO Ray precisa estar instalado. Execute
pip install raypara instalá-lo.table_read = read_builder.new_read() ray_dataset = table_read.to_ray(splits) # Print information about ray_dataset. print(ray_dataset) # Sample output: # MaterializedDataset(num_blocks=1, num_rows=2, schema={key: int64, value: string}) # Print the first two records in ray_dataset. print(ray_dataset.take(2)) # Sample output: # [{'key': 1, 'value': 'AAA'}, {'key': 2, 'value': 'BBB'}] # Convert the entire ray_dataset to a Pandas DataFrame and print the result. print(ray_dataset.to_pandas()) # Sample output: # key value # 0 1 AAA # 1 2 BBB -