解決重複與空值

已完成

資料品質問題如重複記錄與空值,可能會危及你的分析及後續流程。 無論是因同時資料載入產生重複,或因原始資料不完整而出現空值,解決這些問題都是Azure Databricks資料工程師的核心責任。

在這個單元中,你會學習如何利用 SQL 和 PySpark 方法辨識並解決重複、缺失及空值。

識別重複紀錄

在移除重複物品之前,你必須先找到它們。 重複通常發生在同一記錄在同一個資料表中多次出現,通常是因為資料整合問題、重複擷取或同時寫入操作所致。

說明如何辨識重複紀錄的圖表。

使用 GROUP BY 和 HAVING 來尋找重複的項目

辨識重複最直接的方法是將紀錄分組並計數事件。 此 HAVING 子句會過濾分組結果,只顯示出現多次的值。

SELECT 
    customer_id, 
    email, 
    COUNT(*) AS occurrence_count
FROM sales.customers
GROUP BY customer_id, email
HAVING COUNT(*) > 1
ORDER BY occurrence_count DESC;

此查詢會回傳表格中所有 customer_id 和 email 的組合,這些組合多次出現,以幫助你了解重覆的範圍。

在視窗函式中使用 QUALIFY。

此 QUALIFY 子句透過直接過濾視窗函式的結果,提供更優雅的方法,無需子查詢。 這種方法特別適合當你需要看到實際重複的列數,而不只是計數時。

SELECT *
FROM sales.customers
QUALIFY COUNT(*) OVER (PARTITION BY customer_id, email) > 1;

此查詢可能會返回如下結果:

customer_id 電子郵件 名稱 created_at 狀態
1001 jane@contoso.com Jane Doe 2024-01-15 Active
1001 jane@contoso.com 珍·D。 2024-03-20 Active
1002 bob@fabrikam.com 鮑勃·史密斯 2024-02-10 等待中
1002 bob@fabrikam.com 羅伯特·史密斯 2024年2月11日 Active

透過這種模式,你可以從重複的列中取得所有欄位,讓調查重複來源變得更容易。

使用 PySpark 尋找重複的檔案

在 PySpark 中,你可以結合視窗功能與過濾功能來識別重複紀錄:

from pyspark.sql.functions import count
from pyspark.sql.window import Window

window_spec = Window.partitionBy("customer_id", "email")

duplicates_df = (
    df.withColumn("row_count", count("*").over(window_spec))
    .filter("row_count > 1")
    .drop("row_count")
)

display(duplicates_df)

或者,使用此 exceptAll() 方法將原始資料幀與去重版本進行比較:

df_duplicates = df.exceptAll(df.dropDuplicates(["customer_id", "email"]))
display(df_duplicates)

exceptAll() 方法執執行一個保留重複項目的集合差異。 它會返回第一個資料框架中所有不在第二個資料框架中的列,並保留重複出現的情況。 以下是其運作方式:

  1. df.dropDuplicates(["customer_id", "email"]) 建立一個資料框架,每個唯一組合只有一列
  2. exceptAll() 將原始資料框與此去重版本進行比較
  3. 對於每一種獨特的組合,都會配對並刪除一個符合項目,只留下多餘的重複項目

例如,若客戶 1001 在原始資料框中出現三次,在去重版本中出現一次,則 exceptAll() 回傳兩列(多餘的出現次數)。

備註

與 except()不同,後者會移除所有匹配的列並回傳不同的結果,則 exceptAll() 保留了重複的結構。 如果原始資料集中有三個相同的資料列,刪除重復資料集後的資料集中有一個相同的資料列,則 exceptAll() 將傳回兩個資料列。

移除重複紀錄

一旦你辨識出重複,就可以根據你的資料需求選擇解決方案。

說明如何移除重複紀錄的示意圖。

在 PySpark 中移除重複件

此 dropDuplicates() 方法根據指定欄位移除重複列,僅保留第一次出現的欄位:

df_clean = df.dropDuplicates(["customer_id", "email"])

若要在所有欄位中完全刪除列重疊,請使用 distinct():

df_unique = df.distinct()

使用ROW_NUMBER保持特定的行

當你需要控制要保留哪個重複檔案時,可以用 ROW_NUMBER() 和 QUALIFY 依據最新更新等條件來選擇紀錄:

SELECT *
FROM sales.customers
QUALIFY ROW_NUMBER() OVER (
    PARTITION BY customer_id 
    ORDER BY updated_at DESC
) = 1;

此查詢將記錄分割為 customer_id,並依遞減順序排序 updated_at ,並只保留每位客戶最近更新的紀錄。

處理空值與遺失值

空值代表缺失或未知資料。 根據你的分析需求,你可能需要移除有空值的列、以預設值取代,或使用統計補全。

說明如何處理空值和遺失值的圖解。

識別空值

使用 isnull() 函式或 IS NULL 運算子來尋找缺少值的紀錄:

SELECT *
FROM sales.transactions
WHERE amount IS NULL 
   OR customer_id IS NULL;

要計算每欄的空值出現次數:

SELECT 
    COUNT(*) AS total_rows,
    COUNT(*) - COUNT(amount) AS null_amount_count,
    COUNT(*) - COUNT(customer_id) AS null_customer_count
FROM sales.transactions;

刪除無值的列

在 PySpark 中,此 dropna() 方法會移除包含 null 值的列:

# Drop rows where any column is null
df_clean = df.dropna()

# Drop rows where all columns are null
df_clean = df.dropna(how="all")

# Drop rows where specific columns are null
df_clean = df.dropna(subset=["customer_id", "amount"])

參數 how 控制行為: "any" 刪除至少有一個空欄位的列(預設值),而 "all" 當所有指定欄位皆為空時,則只丟棄列。

用預設值填入空值

此 fillna() 方法將空值替換為指定的預設值:

# Fill all null values with a single value
df_filled = df.fillna(0)

# Fill specific columns with different values
df_filled = df.fillna({
    "amount": 0,
    "status": "Unknown",
    "quantity": 1
})

在 SQL 中,使用 COALESCE 函式來提供預設值:

SELECT 
    customer_id,
    COALESCE(amount, 0) AS amount,
    COALESCE(status, 'Unknown') AS status
FROM sales.transactions;

選擇正確的策略

適當的空處理策略取決於你的資料情境:

Scenario 建議的方法
必填欄位 (例如識別碼) 中的 Null 刪除這些列
聚合中的數值欄位的 Null 填入零或平均值
可選類別欄位中的零點 用像「未知」這樣的佔位符填滿
這些空值會扭曲分析 根據領域知識來選擇忽略或填入

小提示

記錄空值處理決策。 未來的分析師需要了解為何某些值會被替換或移除,才能正確解讀查詢結果。

了解如何識別並解決這些資料品質問題,能讓你準備好建立更可靠的資料管線。 有了乾淨的數據,您的下游轉換與分析能產生值得信賴的結果。