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 收尾,那就明天見~
spark/src/main/scala/org/apache/comet/CometConf.scala(shuffle 相關設定與預設值)