この記事では、PyPaimon を使用して Data Lake Formation (DLF) の Paimon テーブルを作成、読み取り、書き込みする方法について説明します。
PyPaimon と DLF の統合
PyPaimon は Apache Paimon の Python SDK です。効率的なデータインジェスト機能を提供し、Python を使用して Paimon テーブルデータを直接読み取り、書き込み、処理できます。
DLF カタログの統合
pypaimon_dlf2 拡張パッケージをインポートして DLF カタログを設定することで、Paimon テーブルのメタデータを Alibaba Cloud Data Lake Formation (DLF) に自動的に同期できます。
この統合により、主に次の利点が得られます。
-
複数エンジン間の相互運用性:メタデータを DLF でホストすることで、MaxCompute、Hologres、Alibaba Cloud EMR などの他の Alibaba Cloud コンピューティングエンジンがこの Paimon データにシームレスにアクセスできます。
-
統合ガバナンス:DLF のデータレイク管理機能により、Paimon テーブルのライフサイクルを管理し、ストレージ形式を自動的に最適化できます。
前提条件
-
DLF データカタログがあること。
-
Python のバージョンは 3.8 以降である必要があります。
python3 --versionを実行して、現在のバージョンを確認してください。
操作手順
ステップ 1:環境の準備
-
次のコマンドを実行して、PyPaimon SDK をインストールします。
pip3 install pypaimon==1.4.1 -
(オプション) インストール後、
pip3 show pypaimonコマンドを実行して、インストールを確認します。
ステップ 2:DLF Paimon テーブルへのアクセス
-
対象ディレクトリで、次のコマンドを実行して
testdlf.pyという名前の新しいファイルを作成します。vim testdlf.py -
testdlf.pyファイルに、次の完全なサンプルコードを追加します。この例では、DLF Paimon テーブルを作成し、読み取りと書き込みを行う方法を示しています。パラメータ設定とデータの読み取り/書き込みの他の方法の詳細については、「コードの詳細」をご参照ください。import pyarrow as pa import pandas as pd from pypaimon import CatalogFactory from pypaimon import Schema # カタログを作成します。 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) # データベースを作成します。 catalog.create_database( name='testdb', ignore_if_exists=True # データベースが既に存在する場合にエラーを無視するかどうかを指定します。 ) # スキーマを作成します。 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' ) # テーブルを作成します。 catalog.create_table( identifier='testdb.tb', schema=schema, ignore_if_exists=True # テーブルが既に存在する場合にエラーを無視するかどうかを指定します。 ) table = catalog.get_table('testdb.tb') # テーブル書き込みとコミット操作を作成します。 write_builder = table.new_batch_write_builder() table_write = write_builder.new_write() table_commit = write_builder.new_commit() # テーブルにデータを書き込みます。PyArrow と Pandas の両方がサポートされています。 # Pandas のサンプルデータを書き込みます。 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) # データをコミットします。 table_commit.commit(table_write.prepare_commit()) # リソースを閉じます。 table_write.close() table_commit.close() # テーブルからデータを読み取ります。複数の出力形式がサポートされています。 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)
ステップ 3:Python ファイルの実行
対象ディレクトリに移動し、次のコマンドで Python スクリプトを実行します。
python3 testdlf.py
次の出力が表示されます。
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: [["08"],["09"]]
key: [[1],[2]]
value: [["AAA"],["BBB"]]
コードの詳細
DLF Paimon テーブルの作成
-
Paimon DLF カタログの作成
説明DLF の Paimon テーブルにアクセスするには、カタログを作成する必要があります。
# catalog_options は、キーと値の両方が文字列である辞書です。 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)次の表に、パラメータの説明を示します。
パラメータ
説明
metastore
固定値
restに設定します。これは、REST Catalog プロトコルを使用して DLF に接続することを示します。dlf.region
DLF リージョンの ID。詳細については、「エンドポイント」をご参照ください。
uri
DLF REST Catalog エンドポイント。VPC 環境では
http://${region_id}-vpc.dlf.aliyuncs.comを使用します。パブリックネットワーク環境ではhttps://dlfnext.${region_id}.aliyuncs.comを使用します。詳細については、「エンドポイント」をご参照ください。warehouse
DLF データカタログの名前。Data Lake Formation コンソールで名前を確認できます。詳細については、「データカタログ」をご参照ください。
dlf.access-key-id
DLF サービスにアクセスするために必要な AccessKey ID。詳細については、「AccessKey ペアの作成」をご参照ください。
dlf.access-key-secret
DLF サービスにアクセスするために必要な AccessKey シークレット。詳細については、「AccessKey ペアの作成」をご参照ください。
token.provider
固定値
dlfに設定します。これは、DLF サービスがアクセストークンを提供することを示します。max-workers
オプション。PyPaimon でデータを読み取る際の同時スレッド数。1 以上の整数である必要があります。デフォルトは 1 で、シリアル読み取りを示します。
-
データベースの作成
Paimon カタログでは、すべてのテーブルは特定のデータベースに属します。テーブルを整理および管理するためにデータベースを作成します。
catalog.create_database( name='database_name', ignore_if_exists=True, # データベースが既に存在する場合にエラーを無視するかどうかを指定します。 properties={'key': 'value'} # オプション。データベースのプロパティ。 ) -
スキーマの作成
スキーマには、列定義、パーティションキー、プライマリキー、テーブルオプション、コメントが含まれます。列定義は
pyarrow.Schemaを使用して記述します。他のパラメータはオプションです。pyarrow.Schemaは、次のいずれかの方法で構築できます。PyArrow
pyarrow.schemaメソッドを使用します。次にコード例を示します。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' )説明pyarrowとPaimon間のデータ型マッピングの詳細については、「PyPaimon データ型マッピング」をご参照ください。Pandas
Pandas データがある場合、
pandas.DataFrameから直接スキーマを導出できます。次にコード例を示します。import pandas as pd import pyarrow as pa from pypaimon import Schema # これはサンプルの DataFrame データです。 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) # DataFrame から pyarrow.Schema を取得します。 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' ) -
テーブルの作成と取得
catalog.create_table( identifier='database_name.table_name', schema=schema, ignore_if_exists=True # テーブルが既に存在する場合にエラーを無視するかどうかを指定します。 ) table = catalog.get_table('database_name.table_name')
テーブルへのデータの書き込み
PyPaimon は現在、bucket オプションが -1 に設定されているプライマリキーテーブルへのデータの書き込みをサポートしていません。
-
テーブル書き込みとコミット操作を作成します。
# テーブル書き込みとコミット操作を作成します。 write_builder = table.new_batch_write_builder() table_write = write_builder.new_write() table_commit = write_builder.new_commit() # Pandas のサンプルデータを書き込みます。 data = { 'date': ['2024-12-01', '2024-12-01', '2024-12-02'], 'hour': ['08', '09', '08'], 'key': [1, 2, 3], 'value': ['AAA', 'BBB', 'CCC'], } -
次のいずれかの方法でテーブルにデータを書き込むことができます。
大規模なデータセットには PyArrow を使用します。通常数ギガバイト以下の小規模なデータセットには、Pandas の方が効率的です。
PyArrow
データは
pyarrow.Tableまたはpyarrow.RecordBatchとして書き込めます。pyarrow.RecordBatchはストリーム処理により適しています。-
方法 1:pyarrow.Table を書き込む
# フィールドを作成します。 fields = [ pa.field('date', pa.string()), pa.field('hour', pa.string()), pa.field('key', pa.int64()), pa.field('value', pa.string()) ] # フィールドからスキーマを作成します。 schema = pa.schema(fields) # テーブルを作成します。 pa_table = pa.Table.from_arrays(data, schema) # データを書き込みます。 table_write.write_arrow(pa_table) -
方法 2:pyarrow.RecordBatch を書き込む
# フィールドを作成します。 fields = [ pa.field('date', pa.string()), pa.field('hour', pa.string()), pa.field('key', pa.int64()), pa.field('value', pa.string()) ] # フィールドからスキーマを作成します。 schema = pa.schema(fields) # RecordBatch を作成します。 record_batch = pa.RecordBatch.from_arrays(data, schema) # データを書き込みます。 table_write.write_arrow_batch(record_batch)
Pandas
pandas.DataFrame からデータを書き込むことができます。
import pandas as pd dataframe = pd.DataFrame(data) table_write.write_pandas(dataframe) -
-
データをコミットし、リソースを閉じます。
# データをコミットします。 table_commit.commit(table_write.prepare_commit()) # リソースを閉じます。 table_write.close() table_commit.close()
テーブルからのデータの読み取り
-
ReadBuilder を作成します。
read_builder = table.new_read_builder() -
PredicateBuilder を使用して、フィルター条件を構築しプッシュダウンします。
-
たとえば、
dateが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) -
たとえば、
keyとvalue列のみを射影できます。read_builder = read_builder.with_projection(['key', 'value'])
説明サポートされているフィルター条件の詳細については、「PyPaimon フィルター条件」をご参照ください。
-
-
splitsを取得します。table_scan = read_builder.new_scan() splits = table_scan.plan().splits() -
splitsをさまざまな出力形式に変換します。Apache Arrow
-
すべてのデータを
pyarrow.Tableに読み込みます。table_read = read_builder.new_read() pa_table = table_read.to_arrow(splits) print(pa_table) # サンプル出力: # pyarrow.Table # key: int64 not null # value: string # ---- # key: [[1],[2]] # value: [["AAA"],["BBB"]] -
データを
pyarrow.RecordBatchReaderに読み込み、バッチを反復処理します。table_read = read_builder.new_read() for batch in table_read.to_arrow_batch_reader(splits): print(batch) # サンプル出力: # pyarrow.RecordBatch # key: int64 # value: string # ---- # key: [1,2] # value: ["AAA","BBB"]
Pandas
データを
pandas.DataFrameに読み込みます。table_read = read_builder.new_read() df = table_read.to_pandas(splits) print(df) # サンプル出力: # key value # 0 1 AAA # 1 2 BBBDuckDB
重要DuckDB をインストールする必要があります。
pip install duckdbを実行してインストールできます。データをインメモリ DuckDB テーブルに変換し、クエリを実行します。
table_read = read_builder.new_read() duckdb_con = table_read.to_duckdb(splits, 'duckdb_table') print(duckdb_con.query("SELECT * FROM duckdb_table").fetchdf()) # サンプル出力: # key value # 0 1 AAA # 1 2 BBB print(duckdb_con.query("SELECT * FROM duckdb_table WHERE key = 1").fetchdf()) # サンプル出力: # key value # 0 1 AAARay
重要Ray をインストールする必要があります。
pip install rayを実行してインストールできます。table_read = read_builder.new_read() ray_dataset = table_read.to_ray(splits) # ray_dataset に関する情報を出力します。 print(ray_dataset) # サンプル出力: # MaterializedDataset(num_blocks=1, num_rows=2, schema={key: int64, value: string}) # ray_dataset の最初の 2 つのレコードを出力します。 print(ray_dataset.take(2)) # サンプル出力: # [{'key': 1, 'value': 'AAA'}, {'key': 2, 'value': 'BBB'}] # ray_dataset 全体を Pandas DataFrame に変換し、結果を出力します。 print(ray_dataset.to_pandas()) # サンプル出力: # key value # 0 1 AAA # 1 2 BBB -