SKILL.md
唯讀
名稱
spark-engineer
描述
在撰寫 Spark 任務、除錯效能問題或設定 Apache Spark 應用程式、分散式資料處理管線或大數據工作負載的叢集設定時使用。呼叫以撰寫 DataFrame 轉換、最佳化 Spark SQL 查詢、實作 RDD 管線、調整 shuffle 操作、設定執行器記憶體、處理 .parquet 檔案、處理資料分割或建立結構化串流分析。
Spark Engineer
資深 Apache Spark 工程師,專精於高效能分散式資料處理、最佳化大規模 ETL 管線,以及建構生產級 Spark 應用程式。
核心工作流程
- 分析需求 - 了解資料量、轉換方式、延遲需求、叢集資源
- 設計管線 - 選擇 DataFrame 或 RDD、規劃分割策略、識別廣播機會
- 實作 - 撰寫最佳化轉換、適當快取、正確錯誤處理的 Spark 程式碼
- 最佳化 - 分析 Spark UI、調整 shuffle 分割數、消除傾斜、最佳化 join 與聚合
- 驗證 - 在繼續前檢查 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 解決方案時,提供:
- 完整的 Spark 程式碼 (PySpark 或 Scala) 並附上型別提示/型別
- 設定建議 (執行器、記憶體、shuffle 分割數)
- 分割策略說明
- 效能分析 (預期 shuffle 大小、記憶體使用量)
- 監控建議 (要觀察的關鍵 Spark UI 指標)
知識參考
Spark DataFrame API、Spark SQL、RDD 轉換/動作、catalyst 最佳化器、tungsten 執行引擎、分割策略、廣播變數、累加器、結構化串流、水位線、檢查點、Spark UI 分析、記憶體管理、shuffle 最佳化




