🗺️ AI 學習與考證地圖
中級科目二程式實戰 · Spark 分散式運算

relativeError 保證名次不保證金額:
聚合交給 Spark,collect 只收小結果

兩段 Spark 程式、兩個考點。 approxQuantile(..., 0.01) 的 0.01 是名次誤差 1%——本站 pyspark 4.2.0 實跑(10 萬筆偏態金額):近似中位數 399 元、精確 404 元——差 5 元,不是 0.01 元,但名次只偏 0.46%(合格)。 top10 管線四站全在 Spark 端分散執行(定義後實測仍是 lazy 的 DataFrame),與 pandas 的 groupby+nlargest 跨引擎逐列一字不差;反面實跑:把全表 collect() 回 driver = 10 萬列 vs 正確的 10 列——差 10,000 倍。口訣:先收斂、再收網

閱讀模式

00題目

資料已是 Spark DataFrame。閱讀下列程式,何者敘述正確?

from pyspark.sql import functions as F   # Spark 的函式工具箱,慣例縮寫 F

median_approx = df.approxQuantile('amount', [0.5], 0.01)[0]   # 近似中位數:名次誤差 ≤ 1%
top10 = (df.groupBy('city')                       # 按城市分組(分散在各節點)
           .agg(F.sum('amount').alias('sales'))   # 各組加總、欄名取作 sales
           .orderBy(F.desc('sales'))              # 按 sales 降冪定序
           .limit(10))                            # 取排序後的前 10 名

先說:總部小辦公室與一萬家分店

Spark 的世界像連鎖超商:driver 是總部的小辦公室executors 是各地分店,資料(發票)分散放在分店裡。選項 B 的世界=總部下令「把每一張發票原件全部寄回總部,總部自己算」——貨車塞爆、辦公室倉庫爆炸(實跑:collect 全表 10 萬列、1.54 秒,而且寄回來之後只是一疊紙——Python list,連 groupBy 都沒得用了)。

正確的玩法是派工單下去、只收結果回來:各分店自己把「本店每城市的銷售額」加總(groupBy+agg)、排好名次(orderBy)、只把前 10 名的摘要寄回總部(limit 後才 collect——10 列)。而 approxQuantile 是另一種聰明:中位數本來要全體排隊才知道,它改用速算草圖回報——但保證書上寫的是「名次誤差不超過 1%」,不是「金額差不超過 0.01 元」。

三個先立好的事實: relativeError=0.01 的保證:回報值的排名落在第 49%~51% 百分位之間(誤差單位是名次,實跑金額差 5 元)。 groupBy/agg/orderBy/limit 都是 transformation——lazy,定義時一列都沒算(實測 type(top10) 仍是 DataFrame);show/collect 這類 action 才開工。 collect 的守則:先讓 Spark 把結果收斂到很小(10 列),再收網拉回 driver——全表 collect 是自廢武功。

先點開看:relativeError 名次誤差lazy 懶惰執行collect 的守則

不熟 Spark 名詞?你需要先認識下列名詞

點擊後出現漸進式說明:白話說明 → 說清楚一點 → 常見錯誤與考點。

分散式的地圖
Spark 是什麼driver 總部executor 分店Spark DataFrame
近似分位數
approxQuantilerelativeError 名次誤差名次 vs 金額速算草圖
top10 的四站
groupBy 分散分組agg 與 aliasorderBy(F.desc)limit(10)
懶惰與收網
lazy 懶惰執行action 扣板機collect 的守則跨引擎對帳

不熟 Python?你需要先認識下列名詞

點擊後出現漸進式說明

帶工具進場
import functions as F'amount' 字串欄名
兩個中括號
[0.5] 分位數清單[0] 取第一個
鏈與回傳
( ) 換行接龍點點接龍回傳 Python 小數

01逐行拆解:一行近似、四站管線

五段,一段一段走完。右上角的「看位置」可以把這一段放回完整程式裡看。

第 1 行帶 F 進場:Spark 的函式工具箱
from pyspark.sql import functions as F   # sum、desc、col……全住在這箱子裡

functions 模組裝著 Spark 的欄位函式——F.sumF.descF.col……慣例縮寫 F(跟 pandas 的 pd、seaborn 的 sns 同一種暗號)。為什麼不用 Python 內建的 sum?因為 F.sum 產生的是「欄位運算式」——一張要發給分店執行的工單,而不是當場把數字加起來。

相關名詞:import functions as FSpark 是什麼

第 2 行approxQuantile:[0.5]、0.01、[0] 三個零件
median_approx = df.approxQuantile('amount', [0.5], 0.01)[0]   # 近似中位數,馬上執行、回傳小數

三個零件:[0.5]分位數清單(0.5=中位數;可以一次要多個 [0.25, 0.5, 0.75]);0.01relativeError——名次誤差的上限:回報值的排名保證落在第 49%~51% 百分位之間;[0] 把回傳清單的第一個值取出來。注意它是 action——這行立刻執行、回傳一個普通的 Python 小數(實跑 399.0)。

A 的死因預告(實跑):精確中位數 404 元、近似 399 元——差 5 元,不是 0.01 元;但名次只偏 0.46%(≤ 1%,合格)。誤差的單位是名次,不是金額——02 節整座實驗室伺候。

相關名詞:approxQuantile[0.5] 分位數清單[0] 取第一個

第 3 行groupBy+agg:分店各自加總
top10 = (df.groupBy('city')                       # 12 個城市分組——工作分散在各節點
           .agg(F.sum('amount').alias('sales'))   # 各組加總;欄名取作 sales

語意跟第 18 頁的 pandas groupby('category')['sales'].sum() 一模一樣,差別在執行的地方:Spark 讓每個節點先加總自己手上的那份資料、再合併小計——資料不動、運算下鄉alias('sales') 幫聚合欄取名——下一站 orderBy 就靠這個名字指認。完全不需要先 collect——分散式聚合正是 Spark 的本業(選項 B 說反了)。

相關名詞:groupBy 分散分組agg 與 alias字串欄名

第 4 行orderBy(F.desc):降冪定序
           .orderBy(F.desc('sales'))   # 按 sales 由大到小——「最高」在這裡定序

F.desc('sales') 明確指定降冪——這行讓下一站的 limit 有了「名次」的意義。對照第 18 頁的口訣:pandas 的 nlargest(10, 'sales') ≡ 這裡的 orderBy(F.desc('sales')).limit(10)——同一題的兩種方言

相關名詞:orderBy(F.desc)跨引擎對帳

第 5 行limit(10):排序後的前十,不是隨機
           .limit(10))   # 取「排好序」的前 10 列——lazy:此刻一列都還沒算

上一行已經按 sales 降冪,limit(10) 取的自然是銷售額前十名——選項 D 的「隨機抽 10 個」在有 orderBy 的管線裡不成立(實跑 top10 嚴格遞減:台北 13,427,368 → … → 屏東 1,275,087,並與 pandas nlargest 逐列一字不差)。更妙的是:這四站全是 transformation——定義完 top10,實測 type(top10) 仍是 DataFrame、一列都沒算——直到 show()/collect() 扣下板機。

一句話記住本題:近似分位數換效率(保證名次不保證金額);聚合排序交給 Spark、collect 只收 10 列小結果——選項 C 的每個字都在這五行裡。

相關名詞:limit(10)lazy 懶惰執行action 扣板機

02誤差的單位:名次,不是金額

選項 A 說 0.01 表示「與精確中位數只差 0.01 元」。本站把 relativeError 轉了四檔(10 萬筆偏態金額,全部實跑):

互動實驗室:誤差的單位relativeError 四檔——金額差幾元 vs 名次偏幾 %
回報的中位數
精確值是 404 元
與精確差幾「元」
A 說永遠 ≤ 0.01 元
名次偏移
保證 ≤ relativeError
本站實跑判決:relativeError=0.01 時回報 399 元、與精確的 404 元差 5 元——「必定只差 0.01 元」當場陣亡;但名次只偏 0.46%(≤ 1%,保證兌現)。轉到 0.1 更戲劇化:回報 318 元、差 86 元——名次偏 9.3% 仍在 10% 保證內。relativeError 的單位是「排名的百分比」:回報值保證落在第 (50−r)%~(50+r)% 百分位之間——金額差幾元,完全看資料分佈長怎樣。

相關名詞:relativeError 名次誤差名次 vs 金額速算草圖

03懶惰的流水線:定義不執行、action 才開工

top10 那條鏈寫完的瞬間,Spark 做了什麼?什麼都沒做——實測給你看:

互動實驗室:懶惰流水線transformation 疊計畫、action 扣板機
lazy 的好處不是偷懶:Spark 把整條管線先看完,才能整體最佳化——它知道你最後只要 10 列,就能把 limit 往前推、讓各節點只回報自己的前十候選,合併後再取總前十——不會真的把全部小計搬來搬去。這就是「宣告做什麼、引擎決定怎麼做」的紅利(跟 SQL 的查詢優化器同一個靈魂)。

相關名詞:lazy 懶惰執行action 扣板機Spark DataFrame

04collect 的代價:B 的世界實跑

選項 B 說「groupBy 前必須先 df.collect()」。把 B 的世界真的跑一遍——順便看看 collect 完你手上剩什麼:

互動實驗室:collect 的代價全表收網 vs 先收斂再收網——實跑對照
B 的三宗罪(實跑): df.collect() 把全表 100,000 列擠回 driver(1.54 秒;資料再大十倍,driver 記憶體直接見底)。 collect 回來的是 Python list——上面根本沒有 groupBy 這個方法,「才能讓 Spark 分散式聚合」在語法層就矛盾:東西都搬離 Spark 了,還分散什麼。 方向整個相反:分散式聚合的正確姿勢是資料不動、運算下鄉——先讓各節點聚合收斂成 10 列,最後才 collect(本題實跑只拉 10 列,相差 10,000 倍)。

相關名詞:collect 的守則driver 總部為什麼要分散式

05跨引擎對帳:Spark 的 top10 = pandas 的 nlargest

top10 到底是不是「銷售額最高的十個城市」?拿第 18 頁的 pandas 寫法對同一份資料重算一次

互動實驗室:跨引擎對帳Spark 管線 vs pandas nlargest——逐列比對
本站實跑 top10(嚴格遞減):台北 13,427,368、新北 12,199,591、台中 9,673,392、高雄 7,971,486、桃園 7,909,057、台南 5,351,797、新竹 3,961,498、彰化 2,030,875、嘉義 1,350,932、屏東 1,275,087——與 pandas groupby('city').sum().nlargest(10) 逐列一字不差(True)。選項 D 的「limit(10) 隨機抽 10 個」不攻自破:orderBy 先定序、limit 才取頭——第 18 頁「先排序 head 才合法」的 Spark 版。

相關名詞:跨引擎對帳limit(10)orderBy(F.desc)

06考場加碼:transformation/action 分類表與 SQL 對照

Spark 題的兩張必備小抄:

類別成員(常考)行為
transformation(疊計畫)groupBy、agg、orderBy、limit、select、filter、withColumn、joinlazy——回傳新 DataFrame,一列不算
action(扣板機)show、collect、count、take、write、approxQuantile觸發整條管線真正執行
-- top10 的 SQL 方言(Spark SQL 也吃這句)
SELECT city, SUM(amount) AS sales   # agg + alias
FROM df                              # Spark DataFrame 註冊成表即可查
GROUP BY city                        # groupBy
ORDER BY sales DESC                  # orderBy(F.desc)
LIMIT 10                             # limit(10)
串起系列:這題是第 18 頁(pandas groupby+nlargest)的大數據方言版——同一個「分組、加總、排名取前 N」的問題,pandas 在單機記憶體算、Spark 派工到叢集算;跨引擎對帳一字不差證明語意相同、執行模型不同。而「先收斂再收網」的紀律,跟第 17 頁「先篩選再 merge」是同一個效率靈魂:能在引擎裡做完的,不要搬出來做。科目二大數據題的核心,考的從來不是語法——是資料該在哪裡被處理

相關名詞:Spark SQL 對照transformation先收斂再收網

07四個選項收工

AapproxQuantile 的 0.01 表示結果必定與精確中位數只差 0.01 元單位就搞錯了

relativeError 的單位是名次(排名百分比)不是金額:0.01 保證回報值的排名落在第 49%~51% 百分位。實跑判決:近似 399 元 vs 精確 404 元——差 5 元(轉到 0.1 檔更差 86 元),但名次只偏 0.46%、完全在保證內。金額差幾元由資料分佈決定——「必定只差 0.01 元」把保證書整個讀錯。

BgroupBy 前必須先 df.collect(),才能讓 Spark 分散式聚合自廢武功

方向整個相反:collect 是把資料搬離 Spark、擠回 driver——實跑全表 100,000 列(1.54 秒),而 collect 回來的是 Python list,上面連 groupBy 都沒有——「才能分散式聚合」在語法層就矛盾。分散式聚合的正確姿勢:資料不動、運算下鄉——groupBy 直接在 DataFrame 上做,collect 留到收斂成 10 列之後(相差 10,000 倍)。

C程式以近似分位數換取效率,聚合與排序仍由 Spark 執行;不應先把全表 collect 到 driver正確

三段敘述三個實跑背書:「近似換效率」——approxQuantile 用速算草圖回報、名次誤差 ≤ 1%(實跑 399 vs 404、名次偏 0.46%);「聚合排序由 Spark 執行」——四站全是 transformation、定義後實測仍是 lazy 的 DataFrame,與 pandas 跨引擎對帳一字不差;「不應全表 collect」——10 列 vs 100,000 列的實跑對照。先收斂、再收網——大數據處理的第一紀律。

Dlimit(10) 會隨機抽 10 個城市,與 sales 排序無關上一行就是排序

limit 的上一站就是 orderBy(F.desc('sales'))——管線按順序執行,limit 取的是排好序的前 10 列。實跑 top10 嚴格遞減(台北 13,427,368 → 屏東 1,275,087),並與 pandas nlargest(10) 逐列一字不差——「隨機」二字無處容身。(沒有 orderBy 的裸 limit 才是「取前面遇到的 N 列」——但那也是「不保證順序」,不是「隨機抽樣」;抽樣是 sample() 的工作。)

回到題目:現在再作答一次

再看一次同一段程式。這次你手上有兩句口訣了:relativeError 保證名次不保證金額;聚合排序交給 Spark、collect 只收小結果

median_approx = df.approxQuantile('amount', [0.5], 0.01)[0]   # 名次誤差 ≤1%(實跑差 5 元)
top10 = (df.groupBy('city').agg(...).orderBy(...).limit(10))   # 四站全在 Spark 端、lazy

08自我檢測

八題,全部都是本題的延伸。答錯會直接告訴你錯在哪。

09重點整理

  1. 兩段程式兩個考點:approxQuantile 一行拿近似中位數(action,馬上回傳小數);top10 四站管線(groupBy → agg → orderBy → limit,全是 lazy 的 transformation)。
  2. 正確答案 C:近似分位數換效率、聚合排序仍由 Spark 執行、不應先把全表 collect 到 driver——三段敘述三個實跑背書。
  3. relativeError 的單位是名次:0.01 保證回報值排名落在第 49%~51% 百分位——不是「差 0.01 元」
  4. A 的死刑(實跑):近似 399 元 vs 精確 404 元——差 5 元、名次偏 0.46%(合格);轉到 0.1 檔差 86 元、名次偏 9.3%(仍在保證內)——金額差幾元由分佈決定。
  5. lazy 實測:定義完 top10,type(top10) 仍是 DataFrame、一列未算;show/collect/count 這類 action 才觸發執行(approxQuantile 也是 action)。
  6. B 的三宗罪(實跑):collect 全表=100,000 列擠回 driver(1.54 秒);collect 回來是 Python list——連 groupBy 都沒有;方向相反——分散式聚合=資料不動、運算下鄉。
  7. collect 守則先收斂、再收網——本題只在 limit(10) 之後拉 10 列(與全表差 10,000 倍)。
  8. D 的死因(實跑):limit 的上一站就是 orderBy(F.desc)——top10 嚴格遞減(台北 13,427,368 → 屏東 1,275,087);隨機抽樣是 sample() 的工作。
  9. 跨引擎對帳(實跑):Spark 管線與 pandas groupby+nlargest(10)(第 18 頁寫法)逐列一字不差——語意相同、執行模型不同。
  10. alias 的角色F.sum('amount').alias('sales') 幫聚合欄命名——orderBy 靠這名字指認(第 16 頁 var_name 的 Spark 版)。
  11. SQL 對照:GROUP BY + ORDER BY DESC + LIMIT 10——Spark SQL 同一句就能跑;transformation/action 分類表要背。
  12. 考場口訣relativeError 保證名次不保證金額;聚合排序交給 Spark、collect 只收小結果——先收斂、再收網
完整程式碼