第3章 データ変換とモデリング / 想定学習時間:30〜40分 / 最終確認:2026年8月

3-4. 重複排除と集約(count・approx_count_distinct・mean・summary)

🎯 この節の学習目標

1. 重複はどこから来るか

重複レコードは、取り込みのリトライ(同じファイルの二重処理)、ソースシステム側の再送、append 書き込みの再実行(3-1)などで自然に発生します。重複を残したまま集計すると売上や件数が水増しされ、結合すれば行が膨張します(3-2)。シルバー層に入れる前の重複排除(deduplication)は、データ品質の基本作業です。

2. dropDuplicates と distinct

メソッド重複の判定基準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行を残したい」という要件には、次のウィンドウ関数を使います。

3. row_number:キーごとに「最新の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 を選びます。

4. 集約の基本:groupBy + agg

ゴールド層(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);

meanavg は同じ関数の別名です。また count("列名")その列が null でない行だけを数え、count("*") は null を含む全行を数える、という違いにも注意してください。

5. 正確性 vs 速度:countDistinct と approx_count_distinct

一意な値の数(カーディナリティ)を数える countDistinct は、正確な代わりに高コストです。一意性を判定するために値を突き合わせる必要があり、大規模データでは大量のシャッフルとメモリを消費します。

そこで用意されているのが approx_count_distinct です。HyperLogLog という確率的アルゴリズムで近似値を高速・省メモリに算出します。誤差は既定で相対標準偏差 5% 程度で、第2引数(rsd)で精度を調整できます(精度を上げるほどコストも上がります)。

countDistinctapprox_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 は正確な値を返す」という選択肢は誤りです。

6. describe と summary:データの統計を素早く把握する

クリーニングや集計の設計前に「データがどんな分布か」を把握するための道具が describesummary です。どちらも探索用の簡易統計を 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、と覚えておけば十分です。なお四分位数は内部的に近似計算であり、ここでも「速度のための近似」という考え方が使われています。

✅ この節のまとめ

練習問題

問1. df.dropDuplicates(["customer_id"]) の動作の説明として正しいものはどれか。

  1. 全列の値が一致する行だけを重複とみなして除去する
  2. customer_id が一致する行を重複とみなし、そのうち最も更新日時が新しい1行を残す
  3. customer_id が一致する行を重複とみなして1行だけ残すが、どの行が残るかは保証されない
  4. customer_id 列そのものを DataFrame から削除する
解答と解説を見る

正解:C

列を指定した dropDuplicates は指定列の一致だけで重複を判定し、残る1行の選ばれ方は保証されません。Aは引数なしの dropDuplicates()(= distinct())の説明です。Bのような「最新を残す」制御をしたければ row_number をウィンドウ関数で使います。Dの列削除は drop() の役割で、まったく別の操作です。

問2. 注文テーブルには同じ order_id の更新履歴が複数行含まれる。updated_at が最も新しい1行だけをキーごとに残す方法として最も適切なものはどれか。

  1. df.dropDuplicates(["order_id"])
  2. Window.partitionBy("order_id").orderBy(col("updated_at").desc()) 上で row_number() を計算し、値が1の行のみ残す
  3. df.distinct()
  4. df.groupBy("order_id").count()
解答と解説を見る

正解:B

キーごとに新しい順で連番を振り rn = 1 を残すのが「最新レコード抽出」の定石です。Aはどの行が残るか不定で、最新である保証がありません。Cは全列一致の重複除去なので、値が少しでも異なる更新履歴はすべて残ってしまいます。Dはキーごとの件数を数えるだけで、元の行は返りません。

問3. 数十億行のアクセスログから日次のユニークユーザー数をダッシュボードに表示したい。数%の誤差は許容できるが、集計時間を大幅に短縮したい。最も適切な関数はどれか。

  1. countDistinct("user_id")
  2. count("user_id")
  3. approx_count_distinct("user_id")
  4. sum("user_id")
解答と解説を見る

正解:C

「誤差許容+高速化」というキーワードは approx_count_distinct(HyperLogLog による近似カウント)の出番です。Aの countDistinct は正確ですが、大規模データではシャッフルとメモリのコストが大きく「時間短縮」の要件に合いません。Bは一意数ではなく非 null 行数を数えるため、同じユーザーの複数アクセスが重複カウントされます。Dは ID の合計というビジネス的に無意味な値です。

問4. 数値列の分布を確認するため、平均・標準偏差に加えて中央値(50%)と四分位数(25%・75%)も一度に見たい。最も適切な方法はどれか。

  1. df.describe() を実行する
  2. df.summary() を実行する
  3. df.count() を実行する
  4. df.printSchema() を実行する
解答と解説を見る

正解:B

summary() は count・mean・stddev・min・max に加えて 25%・50%・75% の四分位数を出力します。Aの describe() は基本5統計のみで四分位数を含みません。Cは行数を1つ返すだけ、Dは列名とデータ型(スキーマ)の表示であり、値の分布は分かりません。