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 --versionuntuk memverifikasi versi saat ini.
Prosedur
Langkah 1: Siapkan lingkungan
-
Jalankan perintah berikut untuk menginstal SDK PyPaimon.
pip3 install pypaimon==1.4.1 -
(Opsional) Setelah instalasi, jalankan perintah
pip3 show pypaimonuntuk memverifikasi instalasi.
Langkah 2: Akses tabel Paimon DLF
-
Di direktori target Anda, jalankan perintah berikut untuk membuat file baru bernama
testdlf.py.vim testdlf.py -
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
-
Buat katalog Paimon DLF.
CatatanAnda 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, gunakanhttps://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.
-
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. ) -
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 membuatpyarrow.Schemadengan 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' )CatatanUntuk informasi tentang pemetaan tipe data antara
pyarrowdanPaimon, 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' ) -
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
PyPaimon saat ini tidak mendukung penulisan data ke tabel kunci primer yang opsi bucket-nya diatur ke -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'], } -
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.Tableataupyarrow.RecordBatch.pyarrow.RecordBatchlebih 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) -
-
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
-
Buat ReadBuilder.
read_builder = table.new_read_builder() -
Gunakan PredicateBuilder untuk membuat dan mendorong kondisi filter.
-
Misalnya, Anda hanya membaca data di mana
dateadalah2024-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
keydanvalue.read_builder = read_builder.with_projection(['key', 'value'])
CatatanUntuk informasi selengkapnya tentang kondisi filter yang didukung, lihat Kondisi filter PyPaimon.
-
-
Ambil
splits.table_scan = read_builder.new_scan() splits = table_scan.plan().splits() -
Konversi
splitske 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.RecordBatchReaderdan 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 BBBDuckDB
PentingAnda harus menginstal DuckDB. Anda dapat menjalankan
pip install duckdbuntuk 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 AAARay
PentingAnda harus menginstal Ray. Anda dapat menjalankan
pip install rayuntuk 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 -