注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。
PySpark の DataFrame は、SQL でテーブルを結合するのと同様に、別のデータフレームや自身と結合できます。データフレームをほかのデータフレームと結合するには、.join() メソッドを使います。このメソッドは、DataFrame、結合に使う列名などの結合条件、および結合方法(left、right、inner など)を受け取ります。
Copied!1df_joined = df_left.join(df_right, 'key', 'left')
df_joined は、df_left.key == df_right.key を条件とした left 結合の結果になります。PySpark は重複する key 列の一方を自動的に削除するため、df_joined には key という名前の列が1列だけ含まれます。
df_left と df_right で結合に使うキーの名前が異なる場合は、結合を実行する前に名前を変更することを推奨します。
明示的な結合条件に使っていないフィールドに同じ名前のものがある場合は、必ず名前を変更するか削除してください。結合が完了すると、これらの名前が衝突します。次のようなループで、DataFrame のすべての列名に特定の接頭辞を付けることができます。
Copied!1 2for column in df.columns: df = df.withColumnRenamed(column, 'some_prefix_' + column)
.join() メソッドは、結合に使うフィールドを単一のフィールドではなくリストとして受け取ることができます。
Copied!1df_joined = df_left.join(df_right, ['column1', 'column2', 'column3'], 'left')
df_joined は、column1、column2、column3 による結合の結果になります。ここでも、df_left と df_right の列名が一致していることを前提としています。
PySpark は、論理演算子を使った任意の式による結合をサポートしています。列 ID が一致し、左側の DataFrame の日付 start が右側の DataFrame の日付 end より前であることを条件に結合する場合を考えます。さらに、特定のフィールド X の内容に応じて、右側の DataFrame の Y に別の値が含まれることを条件に加えるかどうかを決めるものとします。
Copied!1 2 3 4 5 6key_constraint = df_left.ID == df_right.ID date_constraint = df_left.start < df_right.end case_constraint = F.when(df_left.X == 'some_value', df_right.Y == 'some_other_value')\ .otherwise(True) combined_constraints = key_constraint & date_constraint & case_constraint df_joined = df_left.join(df_right, combined_constraints, 'left')
クロス結合を使うと、キーやその他の条件による照合を行わずに、2つのデータフレーム間の行のすべての組み合わせ(直積とも呼ばれます)を生成できます。クロス結合はメモリーやパフォーマンスの問題を引き起こすリスクがあるため、可能な限り避けるべきです。
結果をすぐに絞り込む予定の場合は、クロス結合を使わないでください。代わりに、フィルター条件を結合条件に組み込むと、より効率的に処理できます(上記の詳細な任意の結合条件を参照してください)。
クロス結合を使うには、コードリポジトリで CROSS_JOIN_ENABLED プロファイルを明示的にインポートする必要があります。
Copied!1 2 3 4 5 6 7 8from transforms.api import configure @configure(profile=["CROSS_JOIN_ENABLED"]) @transform_df( ... ) def my_compute_function(input_a, input_b): return input_a.crossJoin(input_b)