Cet article explique comment utiliser PyPaimon pour créer, lire et écrire des données dans les tables Paimon de Data Lake Formation (DLF).
Intégration de PyPaimon et DLF
PyPaimon est le SDK Python d'Apache Paimon. Il offre des capacités d'ingestion de données performantes qui vous permettent de lire, d'écrire et de traiter directement les données des tables Paimon avec Python.
Intégration du catalogue DLF
En important le package d'extension pypaimon_dlf2 et en configurant un catalogue DLF, vous synchronisez automatiquement les métadonnées des tables Paimon vers Alibaba Cloud Data Lake Formation (DLF).
Cette intégration présente les avantages suivants :
Interopérabilité multi-moteurs : Une fois les métadonnées hébergées dans DLF, d'autres moteurs de calcul Alibaba Cloud, tels que MaxCompute, Hologres et Alibaba Cloud EMR, accèdent sans interruption à ces données Paimon.
Gouvernance unifiée : Les fonctionnalités de gestion du data lake de DLF vous permettent de gérer le cycle de vie des tables Paimon et d'optimiser automatiquement leurs formats de stockage.
Prérequis
Vous disposez d'un catalogue de données DLF.
Python version 3.8 ou ultérieure est installé. Exécutez
python3 --versionpour vérifier votre version actuelle.
Procédure
Étape 1 : Préparer l'environnement
-
Exécutez la commande suivante pour installer le SDK PyPaimon.
pip3 install pypaimon==1.4.1 (Facultatif) Après l'installation, exécutez la commande
pip3 show pypaimonpour vérifier l'installation.
Étape 2 : Accéder à une table DLF Paimon
-
Dans votre répertoire cible, exécutez la commande suivante pour créer un nouveau fichier nommé
testdlf.py.vim testdlf.py -
Dans le fichier
testdlf.py, ajoutez l'exemple de code complet suivant. Cet exemple montre comment créer une table DLF Paimon, puis y lire et y écrire des données. Pour plus de détails sur la configuration des paramètres et les autres méthodes de lecture et d'écriture des données, consultez la section Détails du code.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)
Étape 3 : Exécuter le fichier Python
Accédez au répertoire cible et exécutez la commande suivante pour lancer le script Python.
python3 testdlf.py
Le résultat suivant s'affiche.
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"]]
Détails du code
Créer une table DLF Paimon
-
Créer un catalogue Paimon DLF.
RemarqueVous devez créer un catalogue pour accéder aux tables Paimon dans 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)Le tableau suivant décrit les paramètres.
|
**Paramètre**
|
**Description**
| | --- | --- | |
metastore
|
Définissez la valeur fixe `rest`, qui indique que vous vous connectez à DLF via le protocole REST Catalog.
| |
dlf.region
|
L'ID de la région DLF. Pour plus d'informations, consultez la page [Points de terminaison](t2805215.xdita#).
| |
uri
|
Le point de terminaison REST Catalog de DLF. Dans un environnement VPC, utilisez `http://${region_id}-vpc.dlf.aliyuncs.com`. Dans un environnement de réseau public, utilisez `https://dlfnext.${region_id}.aliyuncs.com`. Pour plus d'informations, consultez la page [Points de terminaison](t2805215.xdita#).
| |
warehouse
|
Le nom du catalogue de données DLF. Vous pouvez consulter ce nom dans la console Data Lake Formation. Pour plus d'informations, consultez la section [Catalogue de données](t2803098.xdita#).
| |
dlf.access-key-id
|
L'ID AccessKey requis pour accéder au service DLF. Pour plus d'informations, consultez la section [Créer une paire de clés AccessKey](t162684.xdita#).
| |
dlf.access-key-secret
|
Le secret AccessKey requis pour accéder au service DLF. Pour plus d'informations, consultez la section [Créer une paire de clés AccessKey](t162684.xdita#).
| |
token.provider
|
Définissez la valeur fixe `dlf`, qui indique que le service DLF fournit le jeton d'accès.
| |
max-workers
|
Facultatif. Le nombre de threads simultanés pour la lecture des données dans PyPaimon. Doit être un entier supérieur ou égal à 1. La valeur par défaut est 1, ce qui indique une lecture sérielle.
| -
Créer une base de données.
Dans un catalogue Paimon, chaque table appartient à une base de données spécifique. Créez des bases de données pour organiser et gérer vos tables.
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. ) -
Créer un schéma.
Un schéma comprend les définitions de colonnes, les clés de partition, les clés primaires, les options de table et les commentaires. Les définitions de colonnes sont décrites à l'aide de
pyarrow.Schema. Les autres paramètres sont facultatifs. Vous pouvez construire unpyarrow.Schemade l'une des deux manières suivantes.PyArrow
Utilisez la méthode
pyarrow.schema. L'exemple de code suivant illustre cette approche.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' )RemarquePour connaître le mappage des types de données entre
pyarrowetPaimon, consultez la page Mappage des types de données PyPaimon.Pandas
Si vous disposez de données Pandas, vous pouvez déduire le schéma directement à partir d'un
pandas.DataFrame. L'exemple de code suivant illustre cette approche.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' ) -
Créez et récupérez une table.
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')
Écrire des données dans une table
PyPaimon ne prend pas encore en charge l'écriture de données dans des tables à clé primaire dont l'option bucket est définie sur -1.
-
Créez les opérations d'écriture et de validation (commit) pour la table.
# 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'], } -
Vous pouvez écrire des données dans la table de l'une des manières suivantes :
Pour les grands ensembles de données, privilégiez PyArrow. Pour les ensembles plus petits, généralement inférieurs à quelques gigaoctets, Pandas peut s'avérer plus efficace.
PyArrow
Vous pouvez écrire les données sous forme de
pyarrow.Tableou depyarrow.RecordBatch. pyarrow.RecordBatch convient mieux au traitement en flux continu.-
Méthode 1 : Écrire un 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éthode 2 : Écrire un 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
Vous pouvez écrire des données à partir d'un pandas.DataFrame.
import pandas as pd dataframe = pd.DataFrame(data) table_write.write_pandas(dataframe) -
-
Validez les données et libérez les ressources.
# Commit the data. table_commit.commit(table_write.prepare_commit()) # Close the resources. table_write.close() table_commit.close()
Lire des données depuis une table
-
Créez un ReadBuilder.
read_builder = table.new_read_builder() -
Utilisez un PredicateBuilder pour construire et pousser les conditions de filtrage.
-
Par exemple, vous pouvez lire uniquement les données où
dateest2024-12-01.predicate_builder = read_builder.new_predicate_builder() predicate = predicate_builder.equal('date', '2024-12-01') read_builder = read_builder.with_filter(predicate) -
Par exemple, vous pouvez projeter uniquement les colonnes
keyetvalue.read_builder = read_builder.with_projection(['key', 'value'])
RemarquePour plus d'informations sur les conditions de filtrage prises en charge, consultez la page Conditions de filtrage PyPaimon.
-
-
Récupérez les
splits.table_scan = read_builder.new_scan() splits = table_scan.plan().splits() -
Convertissez les
splitsvers différents formats de sortie.Apache Arrow
-
Lisez toutes les données dans un
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"]] -
Lisez les données dans un
pyarrow.RecordBatchReaderet itérez sur les lots.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
Lisez les données dans un
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
ImportantVous devez installer DuckDB. Vous pouvez exécuter
pip install duckdbpour l'installer.Convertissez les données en une table DuckDB en mémoire et interrogez-la.
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
ImportantVous devez installer Ray. Vous pouvez exécuter
pip install raypour l'installer.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 -