iT邦幫忙

2026 iThome 鐵人賽

DAY 19
0
Software Development

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

Day 19 SortExec 記憶體不夠時怎麼寫磁碟:排序、Pushdown與 Top-K

  • 分享至 

  • xImage
  •  

D18 把 join 演算法選擇拆開來看,最後提一句 Sort、hash 都有 spill 路徑 ,但 spill 到底怎麼運作?
今天把三件事湊在一起說:

  • 排序:資料超過記憶體了,怎麼繼續?
  • 下推:能不能一開始就少讀?
  • 提前終止:Top-K 跟 limit 是不是還有優惠?

看起來是三個題目,其實是同一個原則的三種姿勢,都在問一句話:能不能不做完全部

先看一條會 spill 的 SQL

延續 D17、D18 的例子:

SELECT l_orderkey, SUM(l_extendedprice) AS revenue
FROM lineitem l JOIN orders o ON l.l_orderkey = o.o_orderkey
WHERE o.o_orderdate BETWEEN '1998-01-01' AND '1998-12-31'
GROUP BY l_orderkey
ORDER BY revenue DESC
LIMIT 100

lineitem 是六億筆的事實表,orders 幾千萬筆
SUM 群組完之後 ORDER BY 要排序,機器 heap 假設只有 8 GB,lineitem 掃出來的資料光 raw 就好幾倍, SQL 要跑起來至少三個機會:

  1. o_orderdate BETWEEN ... 這個 filter 可以在 scan 那層就剪掉大部分 orders
  2. SELECT l_orderkeyl_extendedprice,剩下 14 個欄位不用讀
  3. ORDER BY ... LIMIT 100,其實只要維持一個大小為 100 的堆疊就夠

三個機會如果一個都不用,就算演算法選得再漂亮,執行時還是要面對整批資料塞不下記憶體的問題。當那一刻真的到來,就輪到 spill 上場。

排序:外部排序的兩段結構

DataFusion 的 SortExec 是個 pipeline breaker,D16 講過它要看完全部資料才知道最大的是哪一筆。看完全部等於全部要進記憶體,這個假設不成立的時候怎麼辦?

答案是 Graefe 2006 整理過的 external merge sort,兩個階段:

Phase 1,Run generation:一次讀 batch 進來排序,記憶體滿了就把排好的一段 run 寫到磁碟,記憶體釋放繼續讀下一段。跑完之後磁碟上會有 N 個各自排好序的 run。

Phase 2,Multi-way merge:把 N 個 run 打開,每個 run 讀一小塊到記憶體,用一個 min-heap 挑最小的吐出來,補進來繼續挑。整個過程只需要「N 個 batch + 一個 heap」的記憶體,就能產生完整排好序的輸出。

[記憶體]                       [磁碟]

第一段讀進來 → sort → 寫出 →→→ Run 0 [ 排好序 ]
第二段讀進來 → sort → 寫出 →→→ Run 1 [ 排好序 ]
第三段讀進來 → sort → 寫出 →→→ Run 2 [ 排好序 ]
              ...
第 N 段讀進來 → sort → 寫出 →→→ Run N [ 排好序 ]

Phase 2:
Run 0 ┐
Run 1 ┼→ min-heap →→→ 完整排序輸出(串流)
Run 2 ┤
 ...  ┘

這是 60 年代就在磁帶上做的東西,跳到 SSD 世代其實沒變太多,改變的是常數:一個 run 大小從幾百 KB 拉到幾百 MB,merge 從 5-way 拉到幾十 way,寫入放大比從 3x 降到接近 1x。演算法本身還是那個演算法,也就是為什麼今天 DataFusion 的 SortExec 走這條,跟 Postgres 的 tuplesort、Spark 的 UnsafeExternalSorter 走的其實是同一條路。

跟 MemoryPool 的介面回扣 D17:SortExec 是 spillable consumer,跟 pool grow 時被拒絕,就進 spill 路徑寫一個 run 到 tmp,reservation 釋放,pool 額度回來。整個過程 pool 不介入決策,只做批准或拒絕,spill 的方向與時機都是 SortExec 自己決定的。

Pushdown:不要讀那些不會用的資料

排序拆完,回頭看那個 SQL 開頭
lineitem 六億筆,其中 o_orderdate BETWEEN ... 可能只匹配一年的 orders,join 過濾完剩下不到一億筆;SELECT 只用兩個欄位,其他 14 個根本不需要 decode。

這些東西能不能在 scan 那一層就決定,決定了成本的量級。DataFusion 主要有三種下推:

Filter pushdown:把 WHERE 的謂詞下推到 scan。Parquet 這邊搭配 footer 的 min-max 統計,直接跳過整個 row group 不 decode;D18 的 dynamic filter 就是這條路走到極致。

Projection pushdown:只讀 SELECT 用到的欄位。Parquet 是欄式儲存,一個 row group 內每一欄各自一段 column chunk,沒選到的欄位那段 chunk 完全不去碰,I/O 直接省掉。

Limit pushdownORDER BY + LIMIT 這種組合,可以把「我只要前 100 筆」的資訊往下傳到 sort,甚至進到 scan。這一條下面另外講。

下推前後的差別通常不是「快兩倍」,是「掃 100 GB」變成「掃 500 MB」:

不下推                              下推之後

┌─────────────────────────┐        ┌───┐
│ 讀完整 Parquet + 全欄位   │        │ ▓ │  只讀符合 predicate
│ 全部 decode              │        │ ▓ │  的 row group +
│                         │        │ ▓ │  只選要的欄位
│  ~ 100 GB I/O + decode  │        └───┘
│                         │           ~ 500 MB I/O + decode
└─────────────────────────┘

    每一筆都要走完整                    只有留下來的資料
    pipeline                          需要往上走

下推的邏輯位置在 optimizer,不在 executor:optimizer 跑優化 rule 的時候把 Filter / Projection / Limit 這幾個節點「沉」到 TableScan 附近,交給 TableProvider 決定要不要真的接手(用一個「支援哪些 predicate、不支援哪些」的回傳來協商)。這條契約會在 D20 iceberg-rust 那一天正式登場,今天知道位置就好。

Top-K:排序遇上提前終止

回到 ORDER BY revenue DESC LIMIT 100。如果按照 external sort 的定義,要先把整個群組後的資料排完序,才能拿到前 100 筆。可是想一下,真的需要排完嗎?

Top-K 用一個大小為 K 的 heap 就能拿到前 K 筆,時間複雜度 O(N log K),記憶體只要 K 那麼多,不是 O(N log N) 也不是 O(N) 記憶體。DataFusion 的 SortExec 看到下游有 LIMIT K 且 K 相對小時,會走 TopK 這條路:

  • 讀進來的每一 batch,跟目前 heap 裡最差的那一筆比
  • 沒有那筆好就丟掉,好就進 heap 把舊的擠掉
  • 全部讀完,heap 裡就是答案

原本要 sort 六億筆全排完,變成掃過去比 100 筆的 heap,六億次 O(log 100) ≈ 7 次比較。而且完全不需要 spill,因為 heap 那 100 筆一直塞得下。

這是「排序」在遇到「limit」這個資訊之後產生的演算法降級,從外部排序退回到堆疊維護。降級一詞聽起來像慢版本,其實這個「降級」是更快的版本,因為 K 遠小於 N。

三件事的共通形狀:提前終止

回到今天的問題,排序、Pushdown、Top-K 看起來三個獨立題目,但都在做同一件事:在能夠決定這一批不需要的最早時機決定它

  • 下推:能不能在 I/O 開始前就決定不讀(predicate + projection)
  • Top-K:能不能在讀進來後、進 pipeline 前就決定不留(heap 篩掉)
  • Spill:能不能在記憶體塞爆前就決定寫出去(把 memory pressure 換成 disk I/O)

三個看起來一個省 I/O、一個省 memory、一個救命,其實共用一個原則:每一份被留下來繼續往上走的資料,都應該有理由留下來

回到 Comet:shuffle 為什麼是 pipeline breaker,spill 的 fallback 路徑

Shuffle 為什麼是 pipeline breaker 的極端形式:一般 pipeline breaker 是「這個 operator 要看完上游全部才能吐第一筆」,例如 sort、hash join build、aggregate final。Shuffle 更狠,它是跨 process、跨機器的重新分區,所有 mapper 端要把自己這一份資料按 partition 寫到磁碟,reducer 端要等到所有 mapper 都寫完才能開始讀。這個「等」比一般 breaker 多了網路加檔案系統的往返,時間量級再往上跳。所以每一個 shuffle 邊界對 planner 而言都是「非常昂貴」的邊界,能不 shuffle 就不 shuffle,能重用 shuffle 結果就重用。

Spill 的 fallback 路徑是這樣的優先順序:Comet 這邊 native sort 遇到記憶體不夠,第一選擇是自己 spill 到本地 tmp,跟 DataFusion 的路徑一樣。但 Comet 的 MemoryPool 有一條額度線是跟 Spark 的 UnifiedMemoryManager 綁定,Spark 那邊也吃緊時,Comet 這一 batch 會直接標記 unsupported,整個子樹 fallback 給 Spark 的原生 operator 處理,付一次 C2R 的代價。這條 fallback 路徑不是為了效能而設計,是為了保證即使記憶體壓力極大時查詢還能跑完。

Comet 1.0.0 有原生 shuffle 實作,走 Arrow IPC 序列化加上 native 這邊自己管 spill,跟 Spark 原生 shuffle 走完全不同的檔案格式。這條線的細節(分區、序列化、壓縮的三個決策點)留到 D29 shuffle 那一天再拆。

結論

排序超過記憶體時走外部排序,run 生成加 multi-way merge,這是 60 年代的老演算法在 SSD 世代還沒被淘汰。Pushdown 把「能不能不讀」問到 scan 那一層,是最省的省;Top-K 把「能不能不留」問到 sort 那一層,把 O(N log N) 降到 O(N log K)。三件事看起來獨立,共用同一個原則:沒有理由留下來的資料,越早去除越好。Comet 這邊 shuffle 是 pipeline breaker 的極端形式,spill 有從 native 走到 fallback 的降級路徑保命。

明天 D20 進入第四層 iceberg-rust:一個表格式要滿足什麼契約,引擎才能真正把它當資料源?

那就明天見~

參考資料


上一篇
Day 18 Join 演算法:planner 憑什麼選
下一篇
Day 20 Iceberg 怎麼接上引擎:TableProvider、FileScanTask
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言