注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。
Hive スタイルのパーティション分割は、特定の列で絞り込むクエリのパフォーマンスを大幅に向上させるために、データセット内のデータ配置を最適化する方法です。Foundry の Spark ベースのトランスフォームでは、Hive スタイルのパーティション分割は次のように行われます。
Spark や Polars などのコンピューティングエンジンを使用し、これらの列で絞り込むデータセットリーダーは、トランザクションメタデータとファイルパス内のメタデータを自動的に活用して、読み込むファイルを絞り込むことができます。
データ内のパーティション列の値の一意な組み合わせごとに1つ以上のファイルが書き込まれ、過剰な数のファイルを書き込むと、書き込みとその後の読み取りのパフォーマンスが低下します。そのため、Hive スタイルのパーティション分割は、カーディナリティが非常に高い列(一意な値が多数あり、各値に対応する行が少ない列)には適していません。
以下の必要最小限の例では、Python と Java で出力にデータを書き込む際に、Hive スタイルのパーティション分割を設定する方法を示します。
これらの例では、出力に書き込む前に、パーティション列に対して repartitionByRange ↗ を使用してデータフレームを再パーティション分割します。再パーティション分割により、パーティション列の値の一意な組み合わせごとに、_各入力データフレームパーティション内で_1つのファイルが作成されるのではなく、出力全体で1つのファイルのみが作成されるようになります。この再パーティション分割のステップを省略すると、出力データセットに過剰な数のファイルが作成され、書き込みと読み取りのパフォーマンスが低下する可能性があります。
Hive スタイルのパーティション分割では、一般的に repartition ↗ よりも repartitionByRange が推奨されます。これは、repartitionByRange がサンプリングを使用して、データをできるだけ均等に分散するパーティション範囲を推定するためです。一方、repartition はハッシュ関数の結果をパーティション数で割った剰余を使用して、値をデータフレームのパーティションに割り当てます。カーディナリティが低い列では、元のデータが値ごとに比較的均等に分散していても、このハッシュと剰余の演算によってデータが不均等に分散する可能性が高くなります。データの不均等な分散(スキュー)は、Spark エグゼキューターのメモリー不足エラーやジョブの失敗を引き起こす可能性があります。
Copied!1 2 3 4 5 6 7 8 9 10 11 12from transforms.api import transform, Input, Output @transform( transform_output=Output("/path/to/output"), transform_input=Input("/path/to/input"), ) def compute(transform_output, transform_input): transform_output.write_dataframe( transform_input.dataframe().repartitionByRange("record_date", "department"), partition_cols=["record_date", "department"], )
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25package myproject.datasets; import com.palantir.foundry.spark.api.DatasetFormatSettings; import com.palantir.transforms.lang.java.api.Compute; import com.palantir.transforms.lang.java.api.FoundryInput; import com.palantir.transforms.lang.java.api.FoundryOutput; import com.palantir.transforms.lang.java.api.Input; import com.palantir.transforms.lang.java.api.Output; import static org.apache.spark.sql.functions.col; public final class HivePartitioningInJava { @Compute public void myComputeFunction( @Input("ri.foundry.main.dataset.e2dd4bcf-7985-461c-9d08-ee0edd734a1a") FoundryInput myInput, @Output("ri.foundry.main.dataset.4b62bf9b-3700-40f6-9e85-505eaf87e57d") FoundryOutput myOutput) { myOutput.getDataFrameWriter( myInput.asDataFrame().read().repartitionByRange(col("record_date"), col("department"))) .setFormatSettings(DatasetFormatSettings.builder() .addPartitionColumns("record_date", "department") .build()) .write(); } }
repartitionByRange の高度な使用方法上記のコード例では、パーティション数を指定せずに repartitionByRange を呼び出し、Hive スタイルのパーティション分割の設定と同じパーティション列を指定しています。このシンプルな実装は通常は問題ありませんが、非常に大規模なデータを扱う場合、次の2つの状況では問題につながる可能性があります。
repartitionByRange が repartition と同様に、一意な値の組み合わせごとに必ず1つのパーティションに割り当てるためです。spark.sql.shuffle.partitions 設定で構成されたデフォルトのパーティション数を使用します。パーティション列の値の一意な組み合わせの数がその値より大きい場合、1つ以上のパーティションに、複数の値の組み合わせに対応するデータが含まれます。これにより、1つの値の組み合わせに対応するデータが1つの Spark エグゼキューターのメモリーに収まるほど小さい場合でも、Spark のメモリー不足エラーが発生する可能性が高くなります。以下の Python コードサンプルは、複雑さが増す代わりに、これら両方の問題を回避する一般的な実装を示しています。このサンプルのロジックでは、department と record_date の一意な値の組み合わせごとに、データが平均8個の Spark パーティションに分散されます。そのため、出力データセットでは、値の組み合わせごとに1つではなく、およそ8個のファイルが作成されます。
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15from transforms.api import transform, Input, Output @transform( transform_output=Output("/path/to/output"), transform_input=Input("/path/to/input"), ) def compute(transform_output, transform_input): input_df = transform_input.dataframe() unique_date_department_combinations = input_df.select("department", "record_date").distinct().count() partition_count = unique_date_department_combinations * 8 transform_output.write_dataframe( input_df.repartitionByRange(partition_count, "department", "record_date", "record_timestamp"), partition_cols=["department", "record_date"], )