Tous les produits
Search
Centre de documentation

Data Lake Formation:Access DLF using PyPaimon

Dernière mise à jour :Aug 11, 2026

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 --version pour vérifier votre version actuelle.

Procédure

Étape 1 : Préparer l'environnement

  1. Exécutez la commande suivante pour installer le SDK PyPaimon.

    pip3 install pypaimon==1.4.1
  2. (Facultatif) Après l'installation, exécutez la commande pip3 show pypaimon pour vérifier l'installation.

Étape 2 : Accéder à une table DLF Paimon

  1. Dans votre répertoire cible, exécutez la commande suivante pour créer un nouveau fichier nommé testdlf.py.

    vim testdlf.py
  2. 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

  1. Créer un catalogue Paimon DLF.

    Remarque

    Vous 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.
    |



































  2. 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.
    )
  3. 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 un pyarrow.Schema de 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'
    )
    Remarque

    Pour connaître le mappage des types de données entre pyarrow et Paimon, 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'
    )
  4. 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

Remarque

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.

  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'],
    }
  2. 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.Table ou de pyarrow.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)
  3. 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

  1. Créez un ReadBuilder.

    read_builder = table.new_read_builder()
  2. Utilisez un PredicateBuilder pour construire et pousser les conditions de filtrage.

    • Par exemple, vous pouvez lire uniquement les données où date est 2024-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 key et value.

      read_builder = read_builder.with_projection(['key', 'value'])
    Remarque

    Pour plus d'informations sur les conditions de filtrage prises en charge, consultez la page Conditions de filtrage PyPaimon.

  3. Récupérez les splits.

    table_scan = read_builder.new_scan()
    splits = table_scan.plan().splits()
  4. Convertissez les splits vers 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.RecordBatchReader et 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  BBB

    DuckDB

    Important

    Vous devez installer DuckDB. Vous pouvez exécuter pip install duckdb pour 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  AAA

    Ray

    Important

    Vous devez installer Ray. Vous pouvez exécuter pip install ray pour 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