Databricks Certified Data Engineer Associate 教科書
第3章 データ変換とモデリング(Data Transformation and Modeling, 22%)
🎯 この節の学習目標
dropDuplicates(全列/指定列)と distinct、ウィンドウ関数 row_number による最新レコード抽出を使い分けられるgroupBy + agg で count・countDistinct・mean・sum・min/max などの集約を実装できるcountDistinct(正確・低速)と approx_count_distinct(近似・高速)のトレードオフを説明し、describe/summary で統計を把握できる重複レコードは、取り込みのリトライ(同じファイルの二重処理)、ソースシステム側の再送、append 書き込みの再実行(3-1)などで自然に発生します。重複を残したまま集計すると売上や件数が水増しされ、結合すれば行が膨張します(3-2)。シルバー層に入れる前の重複排除(deduplication)は、データ品質の基本作業です。
| メソッド | 重複の判定基準 | SQL 相当 |
|---|---|---|
distinct() | 全列が一致する行 | SELECT DISTINCT * |
dropDuplicates()(引数なし) | 全列が一致する行(distinct と同じ) | SELECT DISTINCT * |
dropDuplicates(["k1","k2"]) | 指定列が一致する行(他の列は違ってもよい) | 直接の相当なし(GROUP BY やウィンドウ関数で代替) |
💡 具体例:全列一致の重複と、キー列だけの重複
df = spark.read.table("bronze.orders_raw")
# (1) 全列が完全に一致する行を1行に(二重取り込みの掃除)
df_unique = df.dropDuplicates() # df.distinct() と同じ
# (2) order_id が同じなら他の列が違っても1行に
df_unique = df.dropDuplicates(["order_id"])
(2)には重要な注意点があります。dropDuplicates(["order_id"]) は同じ order_id の行のうちどの1行が残るかを保証しません。「最新の1行を残したい」という要件には、次のウィンドウ関数を使います。
同じ注文IDに更新履歴が複数行ある場合、「更新時刻が最も新しい行だけ残す」のが典型要件です。ウィンドウ関数 row_number で「キーごとに新しい順で連番を振り、1番だけ残す」と実装します。
💡 具体例:order_id ごとに最新レコードを抽出する
from pyspark.sql import functions as F
from pyspark.sql.window import Window
w = Window.partitionBy("order_id").orderBy(F.col("updated_at").desc())
df_latest = (df
.withColumn("rn", F.row_number().over(w))
.filter(F.col("rn") == 1)
.drop("rn")
)
-- SQL 版
SELECT * EXCEPT (rn)
FROM (
SELECT *,
row_number() OVER (
PARTITION BY order_id
ORDER BY updated_at DESC
) AS rn
FROM bronze.orders_raw
)
WHERE rn = 1;
partitionBy が「どの単位で番号を振り直すか(キー)」、orderBy ... desc() が「何を1番とするか(最新)」を決めます。この形は 3-1 の MERGE 前のソース一意化にもそのまま使えます。
📝 試験のポイント
「キーごとに最新の1件だけ残す」→ row_number() OVER (PARTITION BY キー ORDER BY 時刻 DESC) で rn = 1 を残すが正解の型です。dropDuplicates(["キー"]) はどの行が残るか不定なので、「最新を残す」要件では誤答側に置かれます。類似関数の rank/dense_rank は同順位に同じ番号を振るため、時刻が同一の行が複数残る可能性があり、「必ず1行」にしたいなら row_number を選びます。
ゴールド層(3-6)へ向けた集計の中心が groupBy().agg() です。よく使う集約関数をまとめて確認します。
💡 具体例:地域×月の売上サマリー
from pyspark.sql import functions as F
summary = (spark.read.table("silver.orders")
.withColumn("order_month", F.date_trunc("month", "order_date"))
.groupBy("region", "order_month")
.agg(
F.count("*").alias("order_cnt"), # 行数
F.countDistinct("customer_id").alias("customers"), # 一意顧客数(正確)
F.sum("total_amount").alias("revenue"),
F.avg("total_amount").alias("avg_amount"), # mean でも同じ
F.min("order_date").alias("first_order"),
F.max("order_date").alias("last_order"),
)
)
-- SQL 版
SELECT
region,
date_trunc('month', order_date) AS order_month,
COUNT(*) AS order_cnt,
COUNT(DISTINCT customer_id) AS customers,
SUM(total_amount) AS revenue,
AVG(total_amount) AS avg_amount,
MIN(order_date) AS first_order,
MAX(order_date) AS last_order
FROM silver.orders
GROUP BY region, date_trunc('month', order_date);
mean と avg は同じ関数の別名です。また count("列名") はその列が null でない行だけを数え、count("*") は null を含む全行を数える、という違いにも注意してください。
一意な値の数(カーディナリティ)を数える countDistinct は、正確な代わりに高コストです。一意性を判定するために値を突き合わせる必要があり、大規模データでは大量のシャッフルとメモリを消費します。
そこで用意されているのが approx_count_distinct です。HyperLogLog という確率的アルゴリズムで近似値を高速・省メモリに算出します。誤差は既定で相対標準偏差 5% 程度で、第2引数(rsd)で精度を調整できます(精度を上げるほどコストも上がります)。
countDistinct | approx_count_distinct | |
|---|---|---|
| 結果 | 正確 | 近似(既定で誤差 rsd 5%程度) |
| コスト | 高い(シャッフル・メモリ大) | 低い(高速・省メモリ) |
| 使いどころ | 請求・監査など正確な数が必要な場面 | ダッシュボードの概算、探索的分析、傾向がわかれば十分な場面 |
from pyspark.sql import functions as F
df.agg(
F.countDistinct("user_id").alias("exact_users"), # 正確・低速
F.approx_count_distinct("user_id").alias("approx_users"), # 近似・高速
F.approx_count_distinct("user_id", 0.01).alias("tight"), # 誤差を1%に
).show()
📝 試験のポイント
問題文のキーワードで選びます。「正確な一意数が必要(課金・監査)」→ countDistinct。「巨大データで一意数のおおよその値を高速に知りたい」「多少の誤差は許容できる」→ approx_count_distinct。「approx_count_distinct は正確な値を返す」という選択肢は誤りです。
クリーニングや集計の設計前に「データがどんな分布か」を把握するための道具が describe と summary です。どちらも探索用の簡易統計を DataFrame として返します。
# describe:count / mean / stddev / min / max の5統計
df.select("total_amount", "quantity").describe().show()
# summary:describe の内容+四分位数(25% / 50% / 75%)
df.select("total_amount").summary().show()
# 出したい統計を指定することもできる
df.select("total_amount").summary("count", "min", "50%", "max").show()
| メソッド | 出力される統計 |
|---|---|
describe() | count、mean、stddev、min、max |
summary() | describe の5統計 + 25% / 50% / 75% の四分位数(出力する統計は引数で選択可) |
四分位数まで見たいなら summary、と覚えておけば十分です。なお四分位数は内部的に近似計算であり、ここでも「速度のための近似」という考え方が使われています。
✅ この節のまとめ
distinct() と引数なし dropDuplicates() は全列一致の重複を除去。dropDuplicates(["キー"]) は指定列一致で除去するが、どの行が残るかは不定。row_number() OVER (PARTITION BY キー ORDER BY 時刻 DESC) → rn = 1 で抽出する。groupBy().agg() に count/countDistinct/sum/avg(=mean)/min/max を並べる。count("列") は非 null のみ、count("*") は全行。countDistinct は正確だが高コスト、approx_count_distinct は HyperLogLog による高速な近似(誤差は rsd で調整)。正確性の要否で使い分ける。describe() は基本5統計、summary() はさらに四分位数付き。変換設計の前にデータの分布を把握する。問1. df.dropDuplicates(["customer_id"]) の動作の説明として正しいものはどれか。
customer_id が一致する行を重複とみなし、そのうち最も更新日時が新しい1行を残すcustomer_id が一致する行を重複とみなして1行だけ残すが、どの行が残るかは保証されないcustomer_id 列そのものを DataFrame から削除する正解:C
列を指定した dropDuplicates は指定列の一致だけで重複を判定し、残る1行の選ばれ方は保証されません。Aは引数なしの dropDuplicates()(= distinct())の説明です。Bのような「最新を残す」制御をしたければ row_number をウィンドウ関数で使います。Dの列削除は drop() の役割で、まったく別の操作です。
問2. 注文テーブルには同じ order_id の更新履歴が複数行含まれる。updated_at が最も新しい1行だけをキーごとに残す方法として最も適切なものはどれか。
df.dropDuplicates(["order_id"])Window.partitionBy("order_id").orderBy(col("updated_at").desc()) 上で row_number() を計算し、値が1の行のみ残すdf.distinct()df.groupBy("order_id").count()正解:B
キーごとに新しい順で連番を振り rn = 1 を残すのが「最新レコード抽出」の定石です。Aはどの行が残るか不定で、最新である保証がありません。Cは全列一致の重複除去なので、値が少しでも異なる更新履歴はすべて残ってしまいます。Dはキーごとの件数を数えるだけで、元の行は返りません。
問3. 数十億行のアクセスログから日次のユニークユーザー数をダッシュボードに表示したい。数%の誤差は許容できるが、集計時間を大幅に短縮したい。最も適切な関数はどれか。
countDistinct("user_id")count("user_id")approx_count_distinct("user_id")sum("user_id")正解:C
「誤差許容+高速化」というキーワードは approx_count_distinct(HyperLogLog による近似カウント)の出番です。Aの countDistinct は正確ですが、大規模データではシャッフルとメモリのコストが大きく「時間短縮」の要件に合いません。Bは一意数ではなく非 null 行数を数えるため、同じユーザーの複数アクセスが重複カウントされます。Dは ID の合計というビジネス的に無意味な値です。
問4. 数値列の分布を確認するため、平均・標準偏差に加えて中央値(50%)と四分位数(25%・75%)も一度に見たい。最も適切な方法はどれか。
df.describe() を実行するdf.summary() を実行するdf.count() を実行するdf.printSchema() を実行する正解:B
summary() は count・mean・stddev・min・max に加えて 25%・50%・75% の四分位数を出力します。Aの describe() は基本5統計のみで四分位数を含みません。Cは行数を1つ返すだけ、Dは列名とデータ型(スキーマ)の表示であり、値の分布は分かりません。