iT邦幫忙

2026 iThome 鐵人賽

DAY 8
0
Software Development

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

Day 08 DataFusion SortExec:多欄 ORDER BY 為什麼比單欄貴

  • 分享至 

  • xImage
  •  

嗨嗨~昨天結尾把 EXPLAIN 印出來,plan 最上層那塊 SortExec 還沒沒動,而今天要來看看這塊啦!

問題直白:同一份資料、同一句 ORDER BY,為什麼多加一欄成本就爆走?
「多一欄比較器多一個 memcmp」這解釋不了為什麼換一顆 SET,單欄就能差 4 倍、3 欄能差 50 倍

我們把問題拆成四個具體問句一個一個問:

  • 那顆 SETsort_in_place_threshold_bytes)到底在切什麼?只是 RowFormat 的開關嗎?
  • 同一句 ORDER BY,關掉 in-place 跟打開,差多少?
  • 欄數多一個,成本會怎麼長?線性?平方?
  • 第一欄的基數(打不打得平),會不會直接翻盤?

答案要從一份 500 萬列的 parquet 開始撬。

實驗目的

DataFusion 的 SortExec 裡有一個 if,累積的資料量小於 sort_in_place_threshold_bytes 就整份接起來用欄式比較器排序,超過就每批各自排完再多路合併,只有走合併路徑而且多欄的時候才會把資料編成 RowFormat。我們想看的就是這件事在同一份資料、同一句 SQL 之下差多少,所以拿這個 threshold 當開關,先拉到超大關掉合併路徑,再切回預設打開合併路徑,看 1 欄、2 欄、3 欄各自跑多久,順便再多做一組把第一欄從低基數的 g 換成高基數的 id,看第一欄有沒有打平會不會翻盤。

實驗設備

  • 用 brew 裝的 datafusion-cli 55.0.0
  • 整場實驗把 target_partitions 設成 1,這樣執行計畫裡只會有一個 SortExec,不會出現 SortPreservingMergeExec 這種多路合併節點,讀出來的時間也不會被各核加總放大
  • 資料一份 parquet 放在本機同一顆 SSD 上,路徑固定,不走網路也不走遠端物件儲存

造一份 500 萬列的 parquet

四個欄位:id 每列不同(高基數),g 只取 0 到 7 共 8 個值(低基數,第一欄常打平,故意這樣設計),v 是 double、s 是 md5 字串當 payload。同一份資料就能跑低基數第一欄跟高基數第一欄兩種對照。造資料的 SQL:

COPY (
  SELECT
    value AS id,
    CAST(value % 8 AS INT) AS g,
    CAST(value % 97 AS DOUBLE) AS v,
    md5(CAST(value AS VARCHAR)) AS s
  FROM generate_series(1, 5000000)
) TO 'sortbench.parquet';

generate_series 在 55.0 的欄名是 value 不是 i;PostgreSQL 那套 t(i) / ::INT 這裡不能用,踩過。)

一顆 SET 當開關,兩個值輪流切

要開關 in-place 路徑靠這顆:

-- 關:threshold 拉到 10 GB,強迫走整份接起來的欄式排序
SET datafusion.execution.sort_in_place_threshold_bytes = 10000000000;

-- 開:切回預設 1 MB,強迫走每批各排 + 多路合併
SET datafusion.execution.sort_in_place_threshold_bytes = 1048576;

同時把 target_partitions 鎖成 1,避免 plan 裡冒出 SortPreservingMergeExec 把時間吃掉:

SET datafusion.execution.target_partitions = 1;

每次 SET 完再查一次確認真的切過去了,避免以為關掉其實還開著:

SELECT name, value FROM information_schema.df_settings
 WHERE name = 'datafusion.execution.sort_in_place_threshold_bytes';

六種組合排序 SQL 說明

要驗的 SQL 只有一句,變的只有 ORDER BY 尾巴:

EXPLAIN ANALYZE
SELECT id, g, v, s FROM 'sortbench.parquet'
ORDER BY g, v, s;   -- 或 g / g,v / id,v,s

不能加 LIMIT——加了會走 SortExec: TopK,變成另一題。低基數當第一欄的 ORDER BY gg, vg, v, s 三個 SQL 各跑「關」跟「開」一次,湊出六格;再加碼一組 id, v, s,只把第一欄換成高基數的 id,其他不動。

從 EXPLAIN ANALYZE 讀哪一行

EXPLAIN ANALYZE 印出整棵 plan,我們只盯 SortExec 那一行。3 欄「開」的樣本長這樣:

SortExec: expr=[g@1 ASC NULLS LAST, v@2 ASC NULLS LAST, s@3 ASC NULLS LAST],
  preserve_partitioning=[false],
  metrics=[output_rows=5.00 M, elapsed_compute=909.90ms,
           output_bytes=102.6 GB, output_batches=611,
           spill_count=0, spilled_bytes=0.0 B, spilled_rows=0]

要盯的是 elapsed_compute——SortExec 這個節點花掉的 CPU 時間,target_partitions=1 所以不會被各核加總放大。同時要驗三件事,任一項不對就這格重來:spill_count 必須是 0(不能外溢,不然量到的就不是純排序)、output_rows 必須是 500 萬、整棵 plan 裡不能出現 SortPreservingMergeExecoutput_bytes 在「關」那條路徑會印出上百 GB(同一份 payload 被記到 output_bytes 裡上百次),那是 DataFusion metrics 內部記帳的量,不是真的搬了那麼多位元組,看到別被嚇到。

一鍵重跑

腳本收在同層 refs/experiment/,別台重跑三行就夠:

brew install datafusion
cd refs/experiment
./run.sh            # 造 parquet、暖機、跑六格 + 加碼、印中位數表
./run.sh --parse    # 只重算 raw/ 的中位數

run.sh 對每一格先空跑一次當暖機丟掉、再連跑三次,每次的 SET + EXPLAIN ANALYZE 存進 raw/cell__on|off__rep.sql、原始輸出存進對應 .txt。跑完 parse.py 挑 SortExec 那一行抽 elapsed_compute、順帶檢查 spill_count=0,任一項不對就直接 raise。同時會寫一份 machine.txtdatafusion-cli --version、chip、記憶體、系統版本按下來,之後看數字才知道是哪一台哪一版跑出來的。

實驗結果

暖機丟掉、三次取中位數、單位 ms:

ORDER BY 關(in-place 欄式排序) 開(每批排序 + 多路合併) 開 − 關
g(1 欄,低基數) 28.5 120.3 91.8
g, v(2 欄) 19.2 298.4 279.2
g, v, s(3 欄) 18.1 939.4 921.3
id, v, s(3 欄,高基數第一欄) 17.5 131.4 113.9

三次跑下來 spill_count 都是 0、output_rows 都是 5.00 M、plan 裡也沒有跑出 SortPreservingMergeExec,代表這幾格量到的都是同一個 SortExec 本身的 elapsed_compute,沒有被外溢或多路合併節點汙染。附一個對照,同樣 3 欄同樣的 SQL,「關」那格的 SortExec 那一行只有 elapsed_compute=16.14ms,「開」那格就跳到 909.90ms,這就是那條 921 ms 差額的來源。

代表的意義是什麼

關那一路,欄數幾乎不影響時間——18 到 29 ms 三格的差距落在噪音裡。這條路徑先把所有 batch 接成一整份 RecordBatch,再用 lexsort_to_indices 直接跑欄式比較器,成本被那一次 take 整列 payload 吃掉,多加幾欄比較器只是被最短路徑蓋掉。3 欄 in-place 比 1 欄還略快一點也不用意外,就是次序噪音。

1 欄開跟關就差 4 倍(120 vs 29),但差的其實不是編碼。1 欄的合併路徑在 55.0 有單欄特判,不會編成 RowFormat;差的是 611 路 merge 對上 concat-in-place 的路徑成本本身。這格解釋了「開/關」這個開關不是純 RowFormat 的旋鈕——它同時切了演算法(整份 concat 對每批各排)跟編碼(欄式對 Rows),要孤立編碼要看「開」那一列的欄數變化。

多欄開的成本爆走——120 → 298 → 939,欄數每加一個,時間量級都在往上跳。這就是 RowFormat 編碼加多路合併疊起來付的錢:每一列都要編一份 Rows、每一次 merge 都在做行比較,欄數越多、每次比較搬的位元組越多。這格才是這篇的主軸,也是「多欄比單欄貴這麼多」那個「這麼多」的來源。

第一欄基數決定一切——同樣三欄、同樣走合併路徑,第一欄從低基數的 g 換成高基數的 id,939 掉到 131。memcmp 第一欄就把勝負分光,後面兩欄幾乎不用比。反過來看,g 只有 8 個值,五百萬列平均每 62 萬列就撞一次,第一欄打平之後才會叫醒 vs 兩個比較器。所以「多欄就該用 RowFormat」有前提:第一欄常打平,而且資料量已大到不能 in-place。兩個條件缺一個,這條路徑就從穩賺變成穩虧。

總結

一句話:「多欄比單欄貴這麼多」貴的不是欄本身,貴的是合併路徑——RowFormat 編碼加行比較加多路 merge 疊起來付的錢,欄數只是放大器。in-place 那條路只要塞得下記憶體,欄數怎麼加成本都不動(18–29 ms 三格幾乎持平);「多欄就該用 RowFormat」的教科書結論在這份資料只有第一欄常打平的時候才成立,第一欄一到高基數,同樣三欄同樣走合併,939 立刻掉到 131。

所以下次看到別人建議「多欄排序要開 RowFormat」,可以先反問三件事:資料量多大(塞不塞得下 in-place)、第一欄基數多高(會不會 memcmp 就把勝負分完)、query 有沒有 LIMIT(會不會走 TopK 變成另一題)。三個問題答完,才知道這條建議適不適用。

留一個沒有標準答案的問題:如果第一欄基數決定「多欄開的成本會不會爆走」,這件事該住在哪一層——SortExec 自己在執行時偷看?優化器在規劃階段吃 statistics?還是根本該讓使用者自己 hint?

那就明天見~


上一篇
Day 07 DataFusion 是函式庫,datafusion-cli 是外殼
下一篇
Day 09 Parquet 是什麼、長什麼樣
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言