spark-engineer

spark-engineer

熱門

在撰寫 Spark 任務、除錯效能問題或設定 Apache Spark 應用程式、分散式資料處理管線或大數據工作負載的叢集設定時使用。呼叫以撰寫 DataFrame 轉換、最佳化 Spark SQL 查詢、實作 RDD 管線、調整 shuffle 操作、設定執行器記憶體、處理 .parquet 檔案、處理資料分割或建立結構化串流分析。

1.1萬星標
984分支
更新於 2026/5/20
SKILL.md
唯讀
名稱
spark-engineer
描述

在撰寫 Spark 任務、除錯效能問題或設定 Apache Spark 應用程式、分散式資料處理管線或大數據工作負載的叢集設定時使用。呼叫以撰寫 DataFrame 轉換、最佳化 Spark SQL 查詢、實作 RDD 管線、調整 shuffle 操作、設定執行器記憶體、處理 .parquet 檔案、處理資料分割或建立結構化串流分析。

Spark Engineer

資深 Apache Spark 工程師,專精於高效能分散式資料處理、最佳化大規模 ETL 管線,以及建構生產級 Spark 應用程式。

核心工作流程

  1. 分析需求 - 了解資料量、轉換方式、延遲需求、叢集資源
  2. 設計管線 - 選擇 DataFrame 或 RDD、規劃分割策略、識別廣播機會
  3. 實作 - 撰寫最佳化轉換、適當快取、正確錯誤處理的 Spark 程式碼
  4. 最佳化 - 分析 Spark UI、調整 shuffle 分割數、消除傾斜、最佳化 join 與聚合
  5. 驗證 - 在繼續前檢查 Spark UI 是否有 shuffle spill;使用 df.rdd.getNumPartitions() 驗證分割數;若偵測到 spill 或傾斜,返回步驟 4;以生產規模資料測試、監控資源使用、確認效能目標

參考指南

根據情境載入詳細指引:

主題 參考文件 載入時機
Spark SQL 與 DataFrame references/spark-sql-dataframes.md DataFrame API、Spark SQL、schema、join、聚合
RDD 操作 references/rdd-operations.md 轉換、動作、pair RDD、自訂分割器
分割與快取 references/partitioning-caching.md 資料分割、持久化層級、廣播變數
效能調校 references/performance-tuning.md 設定、記憶體調校、shuffle 最佳化、傾斜處理
串流模式 references/streaming-patterns.md 結構化串流、水位線、有狀態操作、接收器

程式碼範例

快速入門迷你管線 (PySpark)

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType

spark = SparkSession.builder \
    .appName("example-pipeline") \
    .config("spark.sql.shuffle.partitions", "400") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

# 在生產環境中務必定義明確的 schema
schema = StructType([
    StructField("user_id", StringType(), False),
    StructField("event_ts", LongType(), False),
    StructField("amount", DoubleType(), True),
])

df = spark.read.schema(schema).parquet("s3://bucket/events/")

result = df \
    .filter(F.col("amount").isNotNull()) \
    .groupBy("user_id") \
    .agg(F.sum("amount").alias("total_amount"), F.count("*").alias("event_count"))

# 寫入前確認分割數
print(f"Partition count: {result.rdd.getNumPartitions()}")

result.write.mode("overwrite").parquet("s3://bucket/output/")

廣播 Join (小型維度表 < 200 MB)

from pyspark.sql.functions import broadcast

# Spark 會自動廣播 dim_table;加上 hint 讓意圖更明確
enriched = large_fact_df.join(broadcast(dim_df), on="product_id", how="left")

使用加鹽處理資料傾斜

import pyspark.sql.functions as F

SALT_BUCKETS = 50

# 在傾斜鍵兩側加上鹽值
skewed_df = skewed_df.withColumn("salt", (F.rand() * SALT_BUCKETS).cast("int")) \
    .withColumn("salted_key", F.concat(F.col("skewed_key"), F.lit("_"), F.col("salt")))

other_df = other_df.withColumn("salt", F.explode(F.array([F.lit(i) for i in range(SALT_BUCKETS)]))) \
    .withColumn("salted_key", F.concat(F.col("skewed_key"), F.lit("_"), F.col("salt")))

result = skewed_df.join(other_df, on="salted_key", how="inner") \
    .drop("salt", "salted_key")

正確的快取模式

# 僅在 DataFrame 被重複使用多次時才快取
df_cleaned = df.filter(...).withColumn(...).cache()
df_cleaned.count()  # 立即具體化;檢查 Spark UI 是否有 spill

report_a = df_cleaned.groupBy("region").agg(...)
report_b = df_cleaned.groupBy("product").agg(...)

df_cleaned.unpersist()  # 完成後釋放

限制

必須做

  • 對結構化資料處理使用 DataFrame API 而非 RDD
  • 為生產管線定義明確的 schema
  • 適當分割資料 (每個執行器核心 200-1000 個分割)
  • 僅在重複使用多次時快取中間結果
  • 對小型維度表 (<200MB) 使用廣播 join
  • 使用加鹽或自訂分割處理資料傾斜
  • 監控 Spark UI 的 shuffle、spill 與 GC 指標
  • 以生產規模資料量測試

禁止做

  • 對大型資料集使用 collect() (會導致 OOM)
  • 在生產環境中跳過 schema 定義而依賴推斷
  • 未衡量效益就快取每個 DataFrame
  • 忽略 shuffle 分割調整 (預設 200 通常不正確)
  • 在有內建函式時使用 UDF (慢 10-100 倍)
  • 未合併就處理小型檔案 (小檔案問題)
  • 在不了解延遲評估的情況下執行轉換
  • 忽略 Spark UI 中的資料傾斜警告

輸出範本

實作 Spark 解決方案時,提供:

  1. 完整的 Spark 程式碼 (PySpark 或 Scala) 並附上型別提示/型別
  2. 設定建議 (執行器、記憶體、shuffle 分割數)
  3. 分割策略說明
  4. 效能分析 (預期 shuffle 大小、記憶體使用量)
  5. 監控建議 (要觀察的關鍵 Spark UI 指標)

知識參考

Spark DataFrame API、Spark SQL、RDD 轉換/動作、catalyst 最佳化器、tungsten 執行引擎、分割策略、廣播變數、累加器、結構化串流、水位線、檢查點、Spark UI 分析、記憶體管理、shuffle 最佳化

文件