解決重複與空值
資料品質問題如重複記錄與空值,可能會危及你的分析及後續流程。 無論是因同時資料載入產生重複,或因原始資料不完整而出現空值,解決這些問題都是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() 方法執執行一個保留重複項目的集合差異。 它會返回第一個資料框架中所有不在第二個資料框架中的列,並保留重複出現的情況。 以下是其運作方式:
-
df.dropDuplicates(["customer_id", "email"])建立一個資料框架,每個唯一組合只有一列 -
exceptAll()將原始資料框與此去重版本進行比較 - 對於每一種獨特的組合,都會配對並刪除一個符合項目,只留下多餘的重複項目
例如,若客戶 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 | 填入零或平均值 |
| 可選類別欄位中的零點 | 用像「未知」這樣的佔位符填滿 |
| 這些空值會扭曲分析 | 根據領域知識來選擇忽略或填入 |
小提示
記錄空值處理決策。 未來的分析師需要了解為何某些值會被替換或移除,才能正確解讀查詢結果。
了解如何識別並解決這些資料品質問題,能讓你準備好建立更可靠的資料管線。 有了乾淨的數據,您的下游轉換與分析能產生值得信賴的結果。