兩段 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 倍。口訣:先收斂、再收網。
資料已是 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 元」。
點擊後出現漸進式說明:白話說明 → 說清楚一點 → 常見錯誤與考點。
點擊後出現漸進式說明。
五段,一段一段走完。右上角的「看位置」可以把這一段放回完整程式裡看。
from pyspark.sql import functions as F # sum、desc、col……全住在這箱子裡
functions 模組裝著 Spark 的欄位函式——F.sum、F.desc、F.col……慣例縮寫 F(跟 pandas 的 pd、seaborn 的 sns 同一種暗號)。為什麼不用 Python 內建的 sum?因為 F.sum 產生的是「欄位運算式」——一張要發給分店執行的工單,而不是當場把數字加起來。
median_approx = df.approxQuantile('amount', [0.5], 0.01)[0] # 近似中位數,馬上執行、回傳小數
三個零件:[0.5] 是分位數清單(0.5=中位數;可以一次要多個 [0.25, 0.5, 0.75]);0.01 是 relativeError——名次誤差的上限:回報值的排名保證落在第 49%~51% 百分位之間;[0] 把回傳清單的第一個值取出來。注意它是 action——這行立刻執行、回傳一個普通的 Python 小數(實跑 399.0)。
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 說反了)。
.orderBy(F.desc('sales')) # 按 sales 由大到小——「最高」在這裡定序
F.desc('sales') 明確指定降冪——這行讓下一站的 limit 有了「名次」的意義。對照第 18 頁的口訣:pandas 的 nlargest(10, 'sales') ≡ 這裡的 orderBy(F.desc('sales')).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() 扣下板機。
選項 A 說 0.01 表示「與精確中位數只差 0.01 元」。本站把 relativeError 轉了四檔(10 萬筆偏態金額,全部實跑):
top10 那條鏈寫完的瞬間,Spark 做了什麼?什麼都沒做——實測給你看:
選項 B 說「groupBy 前必須先 df.collect()」。把 B 的世界真的跑一遍——順便看看 collect 完你手上剩什麼:
df.collect() 把全表 100,000 列擠回 driver(1.54 秒;資料再大十倍,driver 記憶體直接見底)。② collect 回來的是 Python list——上面根本沒有 groupBy 這個方法,「才能讓 Spark 分散式聚合」在語法層就矛盾:東西都搬離 Spark 了,還分散什麼。③ 方向整個相反:分散式聚合的正確姿勢是資料不動、運算下鄉——先讓各節點聚合收斂成 10 列,最後才 collect(本題實跑只拉 10 列,相差 10,000 倍)。top10 到底是不是「銷售額最高的十個城市」?拿第 18 頁的 pandas 寫法對同一份資料重算一次:
groupby('city').sum().nlargest(10) 逐列一字不差(True)。選項 D 的「limit(10) 隨機抽 10 個」不攻自破:orderBy 先定序、limit 才取頭——第 18 頁「先排序 head 才合法」的 Spark 版。Spark 題的兩張必備小抄:
| 類別 | 成員(常考) | 行為 |
|---|---|---|
| transformation(疊計畫) | groupBy、agg、orderBy、limit、select、filter、withColumn、join | lazy——回傳新 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)
relativeError 的單位是名次(排名百分比)不是金額:0.01 保證回報值的排名落在第 49%~51% 百分位。實跑判決:近似 399 元 vs 精確 404 元——差 5 元(轉到 0.1 檔更差 86 元),但名次只偏 0.46%、完全在保證內。金額差幾元由資料分佈決定——「必定只差 0.01 元」把保證書整個讀錯。
方向整個相反:collect 是把資料搬離 Spark、擠回 driver——實跑全表 100,000 列(1.54 秒),而 collect 回來的是 Python list,上面連 groupBy 都沒有——「才能分散式聚合」在語法層就矛盾。分散式聚合的正確姿勢:資料不動、運算下鄉——groupBy 直接在 DataFrame 上做,collect 留到收斂成 10 列之後(相差 10,000 倍)。
三段敘述三個實跑背書:「近似換效率」——approxQuantile 用速算草圖回報、名次誤差 ≤ 1%(實跑 399 vs 404、名次偏 0.46%);「聚合排序由 Spark 執行」——四站全是 transformation、定義後實測仍是 lazy 的 DataFrame,與 pandas 跨引擎對帳一字不差;「不應全表 collect」——10 列 vs 100,000 列的實跑對照。先收斂、再收網——大數據處理的第一紀律。
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
八題,全部都是本題的延伸。答錯會直接告訴你錯在哪。
type(top10) 仍是 DataFrame、一列未算;show/collect/count 這類 action 才觸發執行(approxQuantile 也是 action)。groupby+nlargest(10)(第 18 頁寫法)逐列一字不差——語意相同、執行模型不同。F.sum('amount').alias('sales') 幫聚合欄命名——orderBy 靠這名字指認(第 16 頁 var_name 的 Spark 版)。「Spark 分散式處理」是中級科目二大數據處理分析與應用的招牌考點,approxQuantile、lazy、collect 三件套輪流出場。最常見的六種問法:① 給 Spark 程式問敘述正誤(本題——驗 relativeError 語意、collect 時機、limit 行為);② 問 relativeError 的意義(名次誤差上限,不是數值誤差);③ 問 transformation 與 action 的分類(lazy 疊計畫 vs 扣板機);④ 問 collect 的正確時機(結果收斂到很小之後);⑤ 問為什麼不把資料拉回單機算(driver 記憶體、網路傳輸、失去平行度);⑥ 問 Spark 與 pandas 的分工(單機放得下用 pandas、放不下派 Spark——語意常常一一對應)。一句口訣:relativeError 保證名次;先收斂、再收網。
approxQuantile 的保證寫在排名上:要求第 q 分位、誤差 r,回報值的真實排名保證落在第 (q−r)×N 到 (q+r)×N 名之間——單位是「名次」。金額差幾元則由資料分佈決定:排名差 1% 在平坦區可能只差幾毛、在陡峭區可能差幾百元。本站實跑(10 萬筆 lognormal):r=0.01 回報 399 元、精確 404 元——差 5 元、名次偏 0.46%(兌現);r=0.1 回報 318 元——差 86 元、名次偏 9.3%(仍兌現)。所以「必定只差 0.01 元」(選項 A)是把保證書的單位整個讀錯——考場看到 relativeError,先默念「這是名次的誤差」。
精確分位數需要「全體排隊」——分散式環境下等於全域排序或多輪重分佈,資料量大時網路與記憶體成本都高。approxQuantile 用 Greenwald-Khanna 這類「速算草圖」:單趟掃描、每個節點只維護一小包帶權重的候選點,最後把小草圖合併——資料不動、只傳草圖。本站 10 萬筆的本機實跑裡近似與精確耗時差不多(都不到一秒)——誠實說:這個規模看不出差距;差距在十億筆、上百節點的叢集上才會現形(草圖幾 KB vs 全域排序)。考點:relativeError=0 也合法——等於要求精確值,代價是走完整計算。
Spark 的 DataFrame 操作分兩類:transformation(groupBy、agg、orderBy、limit、filter、select…)只是「把計畫疊上去」——回傳新的 DataFrame、一列都不算;action(show、collect、count、write、approxQuantile)才扣板機、觸發整條管線執行。實測:定義完 top10,type(top10) 是 DataFrame——什麼 job 都沒跑。lazy 的紅利是整體最佳化:引擎看完全貌才排執行計畫——知道你只要 10 列,就讓各節點先取自己的前十候選再合併,不會傻傻搬運全部小計。這跟 SQL 優化器同一個靈魂:宣告「做什麼」、引擎決定「怎麼做」。
collect 把 DataFrame 的列拉回 driver 變成 Python list。合法時機:結果已經收斂到很小——本題在 limit(10) 之後 collect,只拉 10 列。災難時機:對全表 collect——實跑 100,000 列花 1.54 秒;資料再大幾倍,driver 記憶體直接見底(OOM)。而且 collect 完你手上是 list——沒有 groupBy、沒有 orderBy,Spark 的分散式能力全部歸零(選項 B 的語法矛盾)。替代品:確認資料長相用 show(5)/take(5)(只拉幾列);要轉 pandas 畫圖,先聚合成小表再 toPandas()。守則四個字:先收斂、再收網。
看上一站。有 orderBy 的管線(本題):limit 取「排好序的前 10 列」——實跑 top10 嚴格遞減、與 pandas nlargest 逐列一字不差——這就是「銷售額最高的 10 個城市」。沒有 orderBy 的裸 limit:取「引擎先遇到的 N 列」——省事但不保證任何順序(分區順序決定);注意這叫「未定義順序」,不是「隨機抽樣」——想要統計意義的隨機抽樣要用 sample(fraction=...)。選項 D 把三件事攪在一起:本題有 orderBy(所以是名次)、裸 limit 是未定序(不是隨機)、隨機是 sample 的工作。第 18 頁的口訣搬過來剛好:先排序,limit 才有名次的意義。
逐站對照:groupBy('city') ↔ groupby('city');agg(F.sum('amount').alias('sales')) ↔ ['amount'].sum() 加改名;orderBy(F.desc('sales')).limit(10) ↔ nlargest(10, 'sales')(或 sort_values 降冪+head)。本站對同一份 10 萬筆資料兩邊各跑一次——逐列一字不差(True):語意完全相同,差別只在執行模型(單機記憶體 vs 分散式)。實務選擇:放得進單機記憶體,pandas 快又順手;放不下(或已經在資料湖上),Spark 上場——而你在第 16~18 頁練的 pandas 直覺,翻譯過來九成通用。
三條線。「同一題的方言」:本頁 top10 就是第 18 頁(groupby+nlargest)的 Spark 版——跨引擎對帳一字不差,把「先排序才有名次」的口訣原封搬進大數據世界。「效率的紀律」:「先收斂、再收網」與第 17 頁「先篩選再 merge」同一個靈魂——能在引擎裡做完的,不要搬出來做;資料該在哪裡被處理,是科目二大數據題的核心。「保證書要讀對單位」:relativeError 保證名次不保證金額——跟第 13 頁「0.2857 的分母是誰」、第 18 頁「head 是位置不是名次」同款考點:長得像的數字,先驗明它的單位。