iT邦幫忙

2026 iThome 鐵人賽

DAY 29
0
Software Development

1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~系列 第 29

# Day 29|shuffle:分區、序列化、壓縮,三個通用問題

  • 分享至 

  • xImage
  •  

D11 講過 pipeline-breaking:有些運算子必須把上游全部收完才能吐第一筆。Shuffle 是這件事的極端形式,因為它要收完的東西不只要物化,還要落盤、過網路、在另一台機器上重新組裝。

寫入路徑上有三個決策點。

分區:怎麼把每一列分到目標 partition。Hash-based 最直接,一個 partition 開一個 writer,不用排序;代價是 writer 數量等於 executor cores 乘上 partition 數,記憶體會炸。Comet 的做法是設一個上限(spark.comet.shuffle.jvm.maxWritersPerExecutor,預設 100),超過就退回 sort-based,用一次排序換掉並行 writer 的記憶體。這是典型的「用 CPU 換記憶體」。

序列化:落盤要寫成什麼格式。Spark 原生寫的是 UnsafeRow,列式;Comet 寫的是欄式。欄式的好處有兩個,同型別的值擠在一起壓縮率高很多,而且讀回來直接就是 batch,不用再 R2C 一次。代價是你得自己維護一整套原生 writer 與 reader。

壓縮:Comet 預設 lz4,可換 zstd 或 snappy,zstd 的 level 預設只給 1。這個預設值本身就是答案:shuffle 幾乎總是 I/O bound,壓一點永遠划算,但壓太用力就變成 CPU bound,反而把省下來的時間吐回去。所以大家都停在最便宜的那一格。

三個問題其實是同一個問題的三種問法:這條路上瓶頸在 I/O,所以每一步都在考慮要不要多花一點 CPU 去少搬一點位元組。

今天的思考題:一般的 pipeline-breaking 運算子(像 sort)只是在同一個 task 裡卡住;shuffle 卡住的是整個 stage。這個差別為什麼會讓 shuffle 的優化跟其他運算子長得完全不一樣?

明天 D30 收尾,那就明天見~

參考資料

  • Comet shuffle 說明(user guide)
  • spark/src/main/scala/org/apache/comet/CometConf.scala(shuffle 相關設定與預設值)

上一篇
Day 28 C2R:Comet 算完一定要把資料轉回 Spark 的格式
下一篇
Day 30 驗算 1+1+1>3:加號才是主角
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~30
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言