All Products
Search
Document Center

Data Lake Formation:Akses DLF menggunakan PyPaimon

Last Updated:Jun 19, 2026

Artikel ini menjelaskan cara menggunakan PyPaimon untuk membuat, membaca, dan menulis tabel Paimon di Data Lake Formation (DLF).

Integrasi PyPaimon dan DLF

PyPaimon adalah Python SDK untuk Apache Paimon yang menyediakan kemampuan ingest data efisien, memungkinkan Anda menggunakan Python untuk langsung membaca, menulis, dan memproses data tabel Paimon.

Integrasi katalog DLF

Dengan mengimpor paket ekstensi pypaimon_dlf2 dan mengonfigurasi katalog DLF, metadata tabel Paimon akan secara otomatis disinkronkan ke Alibaba Cloud Data Lake Formation (DLF).

Integrasi ini memberikan manfaat utama berikut:

  • Interoperabilitas multi-engine: Setelah metadata dihosting di DLF, engine komputasi Alibaba Cloud lainnya seperti MaxCompute, Hologres, dan Alibaba Cloud EMR dapat mengakses data Paimon ini secara mulus.

  • Tata kelola terpadu: Kemampuan manajemen data lake DLF memungkinkan Anda mengelola siklus hidup tabel Paimon dan secara otomatis mengoptimalkan format penyimpanannya.

Prasyarat

  • Anda memiliki katalog data DLF.

  • Versi Python harus 3.8 atau lebih baru. Jalankan python3 --version untuk memverifikasi versi saat ini.

Prosedur

Langkah 1: Siapkan lingkungan

  1. Jalankan perintah berikut untuk menginstal SDK PyPaimon.

    pip3 install pypaimon==1.4.1
  2. (Opsional) Setelah instalasi, jalankan perintah pip3 show pypaimon untuk memverifikasi instalasi.

Langkah 2: Akses tabel Paimon DLF

  1. Di direktori target Anda, jalankan perintah berikut untuk membuat file baru bernama testdlf.py.

    vim testdlf.py
  2. Pada file testdlf.py, tambahkan kode contoh lengkap berikut. Contoh ini menunjukkan cara membuat tabel Paimon DLF lalu membaca dan menulis ke dalamnya. Untuk detail konfigurasi parameter dan metode lain dalam membaca serta menulis data, lihat Detail kode.

    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)

Langkah 3: Jalankan file Python

Buka direktori target dan jalankan perintah berikut untuk menjalankan skrip Python.

python3 testdlf.py

Output berikut dikembalikan.

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"]]

Detail kode

Buat tabel Paimon DLF

  1. Buat katalog Paimon DLF.

    Catatan

    Anda harus membuat katalog untuk mengakses tabel Paimon di 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)

    Tabel berikut menjelaskan parameter tersebut.

    Parameter

    Deskripsi

    metastore

    Atur ke nilai tetap rest, yang menunjukkan bahwa Anda terhubung ke DLF menggunakan protokol REST Catalog.

    dlf.region

    ID wilayah DLF. Untuk informasi selengkapnya, lihat Endpoints.

    uri

    Titik akhir DLF REST Catalog. Di lingkungan VPC, gunakan http://${region_id}-vpc.dlf.aliyuncs.com. Di jaringan publik, gunakan https://dlfnext.${region_id}.aliyuncs.com. Untuk informasi selengkapnya, lihat Endpoints.

    warehouse

    Nama katalog data DLF. Anda dapat melihat nama tersebut di Konsol Data Lake Formation. Untuk informasi selengkapnya, lihat Katalog data.

    dlf.access-key-id

    ID AccessKey yang diperlukan untuk mengakses layanan DLF. Untuk informasi selengkapnya, lihat Buat pasangan AccessKey.

    dlf.access-key-secret

    Rahasia AccessKey yang diperlukan untuk mengakses layanan DLF. Untuk informasi selengkapnya, lihat Buat pasangan AccessKey.

    token.provider

    Atur ke nilai tetap dlf, yang menunjukkan bahwa layanan DLF menyediakan token akses.

    max-workers

    Opsional. Jumlah thread konkuren untuk membaca data di PyPaimon. Harus berupa bilangan bulat ≥1. Nilai default adalah 1, yang menunjukkan pembacaan serial.

  2. Buat database.

    Di katalog Paimon, setiap tabel termasuk dalam database tertentu. Buat database untuk mengatur dan mengelola tabel Anda.

    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. Buat skema.

    Skema mencakup definisi kolom, kunci partisi, kunci primer, opsi tabel, dan komentar. Definisi kolom dijelaskan menggunakan pyarrow.Schema. Parameter lain bersifat opsional. Anda dapat membuat pyarrow.Schema dengan salah satu cara berikut.

    PyArrow

    Gunakan metode pyarrow.schema. Kode berikut memberikan contohnya.

    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'
    )
    Catatan

    Untuk informasi tentang pemetaan tipe data antara pyarrow dan Paimon, lihat Pemetaan tipe data PyPaimon.

    Pandas

    Jika Anda memiliki data Pandas, Anda dapat memperoleh skema langsung dari pandas.DataFrame. Kode berikut memberikan contohnya.

    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. Buat dan ambil tabel.

    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')
            

Tulis data ke tabel

Catatan

PyPaimon saat ini tidak mendukung penulisan data ke tabel kunci primer yang opsi bucket-nya diatur ke -1.

  1. Buat operasi tulis dan commit tabel.

    # 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. Anda dapat menulis data ke tabel dengan salah satu cara berikut:

    Untuk dataset besar, gunakan PyArrow. Untuk dataset kecil—biasanya beberapa gigabyte atau kurang—Pandas bisa lebih efisien.

    PyArrow

    Anda dapat menulis data sebagai pyarrow.Table atau pyarrow.RecordBatch. pyarrow.RecordBatch lebih cocok untuk pemrosesan aliran.

    • Metode 1: Tulis 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)
    • Metode 2: Tulis 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

    Anda dapat menulis data dari pandas.DataFrame.

    import pandas as pd
    dataframe = pd.DataFrame(data)
    table_write.write_pandas(dataframe)
  3. Commit data dan tutup sumber daya.

    # Commit the data.
    table_commit.commit(table_write.prepare_commit())
    # Close the resources.
    table_write.close()
    table_commit.close()

Baca data dari tabel

  1. Buat ReadBuilder.

    read_builder = table.new_read_builder()
  2. Gunakan PredicateBuilder untuk membuat dan mendorong kondisi filter.

    • Misalnya, Anda hanya membaca data di mana date adalah 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)
    • Misalnya, Anda hanya memproyeksikan kolom key dan value.

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

    Untuk informasi selengkapnya tentang kondisi filter yang didukung, lihat Kondisi filter PyPaimon.

  3. Ambil splits.

    table_scan = read_builder.new_scan()
    splits = table_scan.plan().splits()
  4. Konversi splits ke berbagai format output.

    Apache Arrow

    • Baca semua data ke dalam 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"]]
    • Baca data ke dalam pyarrow.RecordBatchReader dan iterasi batch-nya.

      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

    Baca data ke dalam 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

    Penting

    Anda harus menginstal DuckDB. Anda dapat menjalankan pip install duckdb untuk menginstalnya.

    Konversi data ke tabel DuckDB dalam memori dan lakukan kueri.

    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

    Penting

    Anda harus menginstal Ray. Anda dapat menjalankan pip install ray untuk menginstalnya.

    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