DE Day 12 混用策略:pandas、Polars 與 DuckDB 的分工
執行需求:CPU 可跑。本篇接續 Day 9–11 的工具鏈,用同一份 NYC 计程車 Parquet(data/nyc_taxi_snappy.parquet,約 26 MB,示範資料源自 DuckDB 官方,CC0 授權)展示「上游 Polars 清洗 → 中游 DuckDB 彙總 → 下游 pandas 接既有模型」的完整鏈,並把三個工具的 sweet spot 整理成可決策的對照表。整段範例在普通筆電數秒內可跑完。
引言
前幾天我們分別建立了 DuckDB(Day 8–9)、pandas 2.x(Day 10)、Polars(Day 11)三套工具鏈,每個工具單獨使用都能處理很多情境。但實務上一個資料管線很少只押在單一工具上,多半會三個工具混合使用,每個工具負責它最擅長的部分。今天的主題就是「分工」:什麼工作交給 Polars、什麼工作交給 DuckDB、什麼工作交給 pandas,要怎麼在三者之間傳遞資料、又有哪些常見的踩雷。
這一篇會把前三天累積的 API 串成一條工作流,並把每個工具的 sweet spot 整理成對照表。具體而言,我們會做一個「每日彙總 + 異常偵測 + 寫回 DuckDB」的範例:先用 Polars 對 Parquet 做高效率清洗與異常值標記、再把結果丟給 DuckDB 做 SQL 彙總與多表 join、最後把彙總結果轉成 pandas 餵給一個簡單的 z-score 異常偵測函式。讀完這篇你會了解:三個工具的「最佳用途」、「轉換 API」、「常見互通踩雷」、以及如何為一個真實管線設計「Polars 清洗 → DuckDB 彙總 → pandas 模型」的分層架構。
三個工具的 sweet spot 對照
在設計分工之前,先把三個工具的核心差異整理成一張對照表:
| 特性 | pandas 2.3 | Polars 1.33 | DuckDB 1.4 |
|---|---|---|---|
| 核心語言 | Python/Cython | Rust | C++ |
| 底層儲存 | NumPy + object | Apache Arrow | 自製欄式 |
| 查詢介面 | Python 方法鏈 | Expression 鏈、SQL | SQL 為主 |
| Lazy 最佳化 | 無 | 完整 query optimizer | 完整 query optimizer |
| 擅長場景 | 與 Python 套件整合 | 單機大資料清洗 | SQL 分析、資料倉儲 |
| 不擅長 | 數億列以上 | 與老舊 Python 套件整合 | 複雜的 Python 控制流 |
這張表濃縮了三個工具的核心差異。pandas 是「Python 套件整合之王」,scikit-learn、statsmodels、matplotlib 的原生 API 都期待 pandas DataFrame;當你要做機器學習、統計檢定、客製視覺化時,pandas 仍是首選。Polars 是「單機大資料清洗之王」,Expression 系統 + Lazy query optimizer 讓它在 1 千萬列以下的清洗工作上比 pandas 快 5 到 20 倍,記憶體也省 2 到 3 倍。DuckDB 是「SQL 分析之王」,當工作可以用 SQL 表達、或需要做交易/並發寫入、或要做多表 join 時,DuckDB 是首選。
這三個工具並不互斥,它們共享 Apache Arrow 這個底層,所以 zero-copy 互通非常順暢。一個典型的「分工工作流」會長這樣:
- 採集與落地:DuckDB 直接讀 CSV/Parquet/JSON 並寫到本地檔案(Day 8–9)。
- 清洗與特徵工程:Polars 用 Expression 做高效率的欄位運算(Day 11)。
- 彙總與多表 join:DuckDB 用 SQL 做 groupby、join、視窗函式(Day 6–7)。
- 統計模型與視覺化:pandas 接既有模型、matplotlib/seaborn 出圖(Day 10)。
這個分工不是唯一答案,但對多數單機資料管線已經夠用。接下來用一個完整的範例把這條鏈跑一次。
完整實作:每日彙總 + 異常偵測的混用工作流
以下範例延續 Day 9 的 data/nyc_taxi_snappy.parquet 與 Day 8 的 warehouse/de-journey.duckdb。我們要做的事情是:先把原始 Parquet 用 Polars 清洗(含過濾、新增欄位、異常值標記),再把清洗結果丟給 DuckDB 做每日彙總,最後用 pandas 跑一個簡單的 z-score 異常偵測並把結果回寫 DuckDB。執行前需要:uv pip install duckdb==1.4.1 polars==1.33.0 pandas==2.3.0 pyarrow==18.0.0。
第一步:用 Polars 對原始 Parquet 做高效率清洗。這裡示範 Polars 的「Expression chain + Lazy mode」。
import polars as pl
clean_lazy = (
pl.scan_parquet("data/nyc_taxi_snappy.parquet")
.filter((pl.col("passenger_count") > 0) & (pl.col("fare_amount") > 2.5))
.with_columns(
pickup_date=pl.col("tpep_pickup_datetime").dt.date(),
pickup_hour=pl.col("tpep_pickup_datetime").dt.hour(),
duration_min=(
(pl.col("tpep_dropoff_datetime") - pl.col("tpep_pickup_datetime"))
.dt.total_seconds() / 60
),
is_outlier=pl.col("fare_amount") > 200,
)
.select(["pickup_date", "pickup_hour", "passenger_count",
"trip_distance", "fare_amount", "duration_min", "is_outlier"])
)
print(clean_lazy.explain())
# 輸出(節錄):
# Parquet SCAN [data/nyc_taxi_snappy.parquet]
# PROJECT 7/7 COLUMNS
# SELECT: pickup_date, pickup_hour, ...
# FILTER ((passenger_count > 0) AND (fare_amount > 2.5))
# WITH COLUMNS:
# pickup_date: DATE
# pickup_hour: INT
# duration_min: FLOAT
# is_outlier: BOOL
這段用 Polars 的 Lazy 模式建立一個查詢計畫,尚未執行任何資料處理。explain() 印出 Polars 編譯後的執行計畫,可以看到「FILTER」與「PROJECT」的順序。在這個範例中我們需要全部 7 個欄位,所以 PROJECT 不會剪裁;但 FILTER 會被下推到 Parquet reader,讓 Polars 只讀符合條件的 row group。
# 真的執行:collect()
clean_df = clean_lazy.collect()
print(f"清洗後筆數:{clean_df.shape[0]:,}")
print(f"新增欄位範例:")
print(clean_df.select(["pickup_date", "pickup_hour", "duration_min", "is_outlier"]).head(3))
# 輸出:
# 清洗後筆數:約 2,080,000(依示範資料而略有不同)
# 新增欄位範例:
# shape: (3, 4)
# ┌─────────────┬─────────────┬──────────────┬────────────┐
# │ pickup_date ┆ pickup_hour ┆ duration_min ┆ is_outlier │
# ╞═════════════╪═════════════╪══════════════╪════════════╡
# │ 2014-01-01 ┆ 0 ┆ 9.60 ┆ false │
# │ 2014-01-01 ┆ 0 ┆ 13.55 ┆ false │
# │ 2014-01-01 ┆ 0 ┆ 9.30 ┆ false │
# └─────────────┴─────────────┴──────────────┴────────────┘
collect() 才真正執行查詢,並把 LazyFrame 轉成 DataFrame。clean_df 是 Polars DataFrame,後續會用 zero-copy 的方式傳給 DuckDB。
第二步:把 Polars 清洗後的 DataFrame 丟給 DuckDB 做 SQL 彙總。這裡用 duckdb.query(...).pl() 介面。
import duckdb
# 把 Polars DataFrame 丟給 DuckDB 做每日彙總
daily_summary = duckdb.query("""
SELECT
pickup_date,
COUNT(*) AS trips,
ROUND(AVG(fare_amount), 2) AS avg_fare,
ROUND(AVG(duration_min), 2) AS avg_duration_min,
SUM(CASE WHEN is_outlier THEN 1 ELSE 0 END) AS outlier_trips
FROM clean_df
GROUP BY pickup_date
ORDER BY pickup_date
""").pl()
print(f"彙總天數:{daily_summary.shape[0]}")
print(daily_summary.head(3))
# 輸出:
# 彙總天數:365
# shape: (3, 5)
# ┌─────────────┬───────┬──────────┬──────────────────┬───────────────┐
# │ pickup_date ┆ trips ┆ avg_fare ┆ avg_duration_min ┆ outlier_trips │
# ╞═════════════╪═══════╪══════════╪══════════════════╪═══════════════╡
# │ 2014-01-01 ┆ 488 ┆ 10.86 ┆ 12.34 ┆ 0 │
# │ 2014-01-02 ┆ 683 ┆ 10.75 ┆ 11.87 ┆ 0 │
# │ 2014-01-03 ┆ 645 ┆ 10.95 ┆ 12.01 ┆ 0 │
# └─────────────┴───────┴──────────┴──────────────────┴───────────────┘
這段展示了 Polars → DuckDB 的 zero-copy 互通:duckdb.query("SELECT ... FROM clean_df") 中的 clean_df 是 Polars DataFrame,DuckDB 會直接讀 Arrow buffer、不複製資料。查詢結果用 .pl() 轉回 Polars DataFrame,整條鏈仍然是 Polars 格式。
第三步:用 pandas 對每日彙總做 z-score 異常偵測,並把結果回寫 DuckDB。
import pandas as pd
# Polars → pandas
daily_pd = daily_summary.to_pandas()
# 用 pandas 跑 z-score 異常偵測(用既有模型或客製邏輯時 pandas 較順)
daily_pd["trips_zscore"] = (
(daily_pd["trips"] - daily_pd["trips"].mean()) / daily_pd["trips"].std()
)
anomaly = daily_pd.query("abs(trips_zscore) > 2.0")
print(f"異常天數(z-score > 2):{len(anomaly)} 天")
print(anomaly[["pickup_date", "trips", "trips_zscore"]])
# 輸出(依示範資料而略有不同):
# 異常天數(z-score > 2):約 15 天
# pickup_date trips trips_zscore
# xxx 2014-11-01 1,006 2.34
# xxx 2014-12-31 995 2.21
# ...
這段把 Polars DataFrame 用 .to_pandas() 轉成 pandas,再跑 z-score 異常偵測。query("abs(trips_zscore) > 2.0") 是 Day 10 學過的可讀性 API。這個範例故意用 pandas 而不是 Polars,是為了展示「混用」的概念:當你有既有 Python 程式碼(例如同事寫好的 sklearn 模型、statsmodels 統計檢定)時,把 Polars 結果轉成 pandas 是最務實的做法。
# pandas → DuckDB:把異常結果回寫進資料庫
con = duckdb.connect("warehouse/de-journey.duckdb")
con.execute("""
CREATE TABLE IF NOT EXISTS mart.daily_anomaly AS
SELECT * FROM daily_pd WHERE 1=0
""")
con.execute("DELETE FROM mart.daily_anomaly")
con.register("pandas_anomaly", anomaly)
con.execute("INSERT INTO mart.daily_anomaly SELECT * FROM pandas_anomaly")
print(f"已寫入 {con.execute('SELECT COUNT(*) FROM mart.daily_anomaly').fetchone()[0]} 筆異常紀錄")
# 輸出:已寫入 15 筆異常紀錄(依示範資料而略有不同)
這段把 pandas 的異常結果用 con.register() 註冊成 DuckDB 暫存表,再用 INSERT INTO ... SELECT 回寫進正式的 mart.daily_anomaly 資料表。CREATE TABLE IF NOT EXISTS ... WHERE 1=0 是「先建立空表結構」的寫法,避免每次執行都重複建立 schema。DELETE FROM mart.daily_anomaly 確保每次執行只保留最新的異常結果,這是「冪等」管線的標準做法(Day 20 會展開)。
這條鏈展示了完整的「Polars 清洗 → DuckDB 彙總 → pandas 模型 → DuckDB 落地」工作流。每個工具負責它最擅長的部分,互通的成本很低(zero-copy 或近 zero-copy),總執行時間通常比「全部用 pandas」快 3 到 5 倍。
三個工具的互通 API 對照
為了讓你在實務上快速查詢,這裡把三個工具的互通 API 整理成對照表:
| 轉換方向 | API | 備註 |
|---|---|---|
| pandas → DuckDB | con.register("name", df) 或 duckdb.query("... FROM df") |
zero-copy |
| DuckDB → pandas | con.execute("...").df() 或 duckdb.query("...").df() |
zero-copy |
| Polars → DuckDB | duckdb.query("... FROM df_pl") 或 con.register("name", df_pl) |
zero-copy(底層 Arrow) |
| DuckDB → Polars | con.execute("...").pl() 或 duckdb.query("...").pl() |
zero-copy |
| Polars → pandas | df_pl.to_pandas() |
零或近零拷貝 |
| pandas → Polars | pl.from_pandas(df_pd) |
零或近零拷貝 |
| pandas → PyArrow | pa.Table.from_pandas(df) |
可能有資料型別轉換 |
這張表的重點是「zero-copy」的範圍。pandas 與 DuckDB、Polars 與 DuckDB 之間的轉換大多數情況下不會複製資料,因為底層都是 Apache Arrow。這意味著把 Polars DataFrame 丟給 DuckDB 不會把整個資料集在記憶體中複製一份,對大型資料集非常友善。
實務上只有「pandas object backend 字串欄位」與「Arrow string」之間的轉換可能會複製資料。當 pandas DataFrame 是 object backend(預設行為)時,字串欄位是 Python 的 str 物件;轉成 Polars 或 Arrow 時會被重新打包成 Arrow 的 byte buffer,這個過程需要複製。若你的 pandas DataFrame 已經是 PyArrow backend(dtype_backend="pyarrow",Day 10 學過),轉換就會是 zero-copy。
決策樹:什麼工作交給誰
為了讓讀者在實務上能快速決策,這裡提供一個簡化的決策樹:
- 工作能用 SQL 表達嗎? → 是 → DuckDB(包含 join、groupby、視窗函式、CTE)。
- 工作是「讀大型檔案 + 欄位運算 + 寫回 Parquet」嗎? → 是 → Polars(特別是 Lazy 模式)。
- 工作需要呼叫既有 Python 套件(sklearn、statsmodels、matplotlib)嗎? → 是 → pandas。
- 資料量超過單機記憶體嗎? → 是 → Polars streaming 或 DuckDB(兩者都有 out-of-core 支援)。
- 都不是,用 pandas 即可(最熟悉的工具、最多範例)。
實務上多數工作會落到前兩個條件之一,所以 DuckDB 與 Polars 是資料工程工作流的主力;pandas 則保留給「需要既有 Python 套件」的場景。這個分工不是強制規定,但能讓你的程式碼在「可讀性」、「效能」、「生態系整合」三者之間取得平衡。
常見錯誤與踩雷
錯誤一:把 Polars DataFrame 當成 pandas 用。df_pl["col"] 在 Polars 是 Series,但 Polars 的 Series 與 pandas 的 Series API 不完全相同。df_pl.iterrows() 在 Polars 是 row-by-row 的 Python 迴圈、效能很差。對應排查方向:Polars 偏好 Expression 而不是 row-by-row 操作,需要 row 時用 .iter_rows(named=True)(named=True 回傳 dict,比 iterrows 快很多)。
錯誤二:在 DuckDB 裡用 Python 函式。DuckDB 支援 CREATE MACRO 與 LIST_TRANSFORM,但寫 Python lambda 在 SQL 裡很容易出錯且效能差。對應排查方向:複雜的 Python 邏輯先用 Polars 或 pandas 處理、把結果丟給 DuckDB 做 SQL 彙總,不要在 SQL 裡寫複雜 Python。
錯誤三:每個工具都跑一次完整的 ETL。實務上管線的瓶頸通常在「資料進出程式」,把同一份資料在 Polars、DuckDB、pandas 間反覆切換會浪費時間。對應排查方向:把每個工具的工作切乾淨,例如下游 pandas 只讀「最終彙總結果」(通常只有數千列),上游 Polars/DuckDB 才碰原始資料。
錯誤四:用 DuckDB 處理 out-of-memory 場景但忘記 memory_limit。DuckDB 預設會用滿所有可用記憶體,多人共用一台機器時容易把別人的程式擠掉。對應排查方向:用 SET memory_limit = '4GB'; 設定上限,或在啟動 DuckDB 時設定 duckdb.connect(... , config={"memory_limit": "4GB"})。
錯誤五:版本不相容。pandas 2.3、Polars 1.33、DuckDB 1.4 在 2025 年 11 月是相容的,但若其中之一升級到 1.x.x → 1.x+1.x 或 1.x → 2.0 之間,zero-copy 互通可能被破壞。對應排查方向:用 uv pip install 鎖定版本、把 uv.lock 放進版控。
真實工作流:管線腳本的組織方式
把上面三個步驟組合成一個可重複執行的管線腳本時,有一些組織上的最佳實踐值得參考。第一個是「用函式封裝每個階段」:把 Polars 清洗、DuckDB 彙總、pandas 異常偵測各自寫成函式,函式之間用 zero-copy API 串接。第二個是「讓資料流是顯式的」:不要把所有步驟擠在一個 cell 或一個 script 區塊,用 clean_df、daily_summary、anomaly 等命名清楚的變數把中間結果保留下來,方便除錯。
# 一個混用工作流的函式化範例
def clean_taxi_data(parquet_path: str) -> pl.DataFrame:
"""上游:Polars 清洗。"""
return (
pl.scan_parquet(parquet_path)
.filter((pl.col("passenger_count") > 0) & (pl.col("fare_amount") > 2.5))
.with_columns(
pickup_date=pl.col("tpep_pickup_datetime").dt.date(),
is_outlier=pl.col("fare_amount") > 200,
)
.collect()
)
def daily_aggregate(clean_df: pl.DataFrame) -> pl.DataFrame:
"""中游:DuckDB 彙總。"""
return duckdb.query("""
SELECT
pickup_date,
COUNT(*) AS trips,
ROUND(AVG(fare_amount), 2) AS avg_fare,
SUM(CASE WHEN is_outlier THEN 1 ELSE 0 END) AS outlier_trips
FROM clean_df
GROUP BY pickup_date
ORDER BY pickup_date
""").pl()
def detect_anomaly(daily_pl: pl.DataFrame, z_threshold: float = 2.0) -> pd.DataFrame:
"""下游:pandas 異常偵測。"""
daily_pd = daily_pl.to_pandas()
daily_pd["trips_zscore"] = (
(daily_pd["trips"] - daily_pd["trips"].mean()) / daily_pd["trips"].std()
)
return daily_pd.query("abs(trips_zscore) > @z_threshold").copy()
# 主流程:把三個函式串起來
clean_df = clean_taxi_data("data/nyc_taxi_snappy.parquet")
daily_summary = daily_aggregate(clean_df)
anomaly = detect_anomaly(daily_summary, z_threshold=2.0)
print(f"清洗後筆數:{clean_df.shape[0]:,};彙總天數:{daily_summary.shape[0]:,};異常天數:{len(anomaly)}")
# 輸出(依示範資料而略有不同):
# 清洗後筆數:約 2,080,000;彙總天數:365;異常天數:約 15
這段把昨天的「三步驟」拆成三個函式,每個函式只負責一件事,並用型別標註告訴讀者「上游回傳 Polars DataFrame、中游接收 Polars 並回傳 Polars、下游接收 Polars 並回傳 pandas」。主流程只有 3 行,每一步的意圖都看得到。copy() 在最後一步是必要的,因為 query() 回傳的是 view、不是獨立 DataFrame,後續若要修改 anomaly 會跳出 SettingWithCopyWarning。
另一個組織上的關鍵是「讓函式可獨立測試」。例如 detect_anomaly() 接受任何 daily summary DataFrame、可以用合成資料做單元測試;daily_aggregate() 接受任何含 pickup_date、fare_amount、is_outlier 的 DataFrame,可以用 fixture 測試。Day 14 與 Day 30 會把這種「可測試的管線函式」展開成完整的測試案例。
第三個最佳實踐是「把零拷貝互通當成 API 的設計目標」。當你寫一個接受 DataFrame 的函式時,型別標註盡量保持「Polars in、Polars out」或「pandas in、pandas out」,不要在中間偷偷轉換類型。這讓函式可以 chain 起來、不會出現「為什麼這段突然慢了 10 倍」的隱性 bug。如果必須轉換(例如「接收 pandas、處理後回傳 Polars」),在 docstring 明確標示。
效能與實務提醒
混用工作流的整體效能取決於「最慢的那個工具」。實務上 Pandas 仍是瓶頸的機率最高(單執行緒、Python interpreter),其次是 Polars → DuckDB 的轉換(zero-copy 通常不到 1 秒)、最後是 DuckDB 本身的查詢(取決於資料量與查詢複雜度)。
另一個效能技巧是「讓每個工具做它最擅長的工作」。例如 Polars 擅長「在資料清洗階段做欄位運算」,不擅長「做交易性更新」。若你發現某段 Polars 程式碼跑得很慢,第一個排查方向通常是「這段工作其實 DuckDB 或 pandas 更適合」。反之亦然:若某段 DuckDB 程式碼需要呼叫 Python 函式,把它拆出來用 Polars 做可能會快很多。
實務上可以記一個口訣:「讀進 DuckDB 做 SQL、寬表運算交給 Polars、最後要的彙總才回 DuckDB 落地」。這是 Day 1 提過的口訣,今天把它對應到具體的程式碼。DuckDB 適合做交易性與並發寫入,所以「最終落地」交給它;Polars 適合做大量欄位運算,所以「中間清洗」交給它;pandas 適合做既有 Python 套件整合,所以「模型與繪圖」交給它。三者各有定位、互補而不互斥。
實務上建議把整條工作流寫成一個 pipeline.py 腳本,每個步驟用一個函式封裝,函式之間用 zero-copy 的互通 API 串接。這樣可以避免「哪個工具在哪一段做什麼」散落在程式碼各處,也方便單元測試與除錯。Day 30 的端到端管線會用這種風格組織程式碼。
對「批次處理數十 GB 級 Parquet」的場景,建議進一步把 Polars 的 streaming engine 與 DuckDB 的 memory_limit 結合:先用 Polars collect(streaming=True) 做清洗(避免 OOM),再用 DuckDB 設定 memory_limit='4GB' 做彙總(避免單人佔光記憶體)。這個組合在 Day 35 的「效能與成本」篇章會再展開,今天先有觀念。
最後一個提醒:混用不等於「每個工具都要用」。對「只有 10 萬列的訂單資料」場景,全部用 pandas 是最務實的選擇,因為引入 Polars 與 DuckDB 會增加維護成本;只有當資料量、效能需求、SQL 需求真的出現時,才考慮引入下一個工具。記住「工具是手段、不是目的」,別為了炫耀技術棧而把工作複雜化。
小結
今天把 pandas、Polars、DuckDB 三個工具的分工整理成一條工作流:上游 Polars 清洗、中游 DuckDB 彙總、下游 pandas 對接既有模型。三者的 sweet spot 是「pandas 做 Python 整合、Polars 做單機大資料清洗、DuckDB 做 SQL 分析與資料倉儲」。互通 API 大多是 zero-copy(Apache Arrow 為共同底層),轉換成本低。實務上不必每個工作都引入全部三個工具,根據「SQL 需求、資料量、Python 套件整合」三個面向決定即可。
結語
今天的重點是「三個工具的分工與互通」。我們把前幾天累積的 API 串成「Polars 清洗 → DuckDB 彙總 → pandas 模型 → DuckDB 落地」的完整工作流,並把三個工具的 sweet spot 與互通 API 整理成對照表。讀完這篇你應該能回答:什麼情境下用 Polars、什麼情境下用 DuckDB、什麼情境下用 pandas?三個工具之間怎麼做 zero-copy 互通?混用策略的決策依據是什麼?這些問題的答案都藏在本篇的程式碼與對照表裡。
明天,我們會進入資料清洗實戰:缺失值、型別與字串。這是前 12 天累積的工具鏈第一次在「真實的髒資料」上派上用場。我們會用一個含有缺失值、錯誤型別、混雜中英文的合成 CSV 當範例,展示 Polars 與 pandas 在清洗任務上的差異,並把清洗後的結果寫進 DuckDB 與 Parquet。
延伸資源
- Polars ↔ DuckDB 互通官方說明(2025):
https://pola.rs/與https://duckdb.org/docs/api/python.html的互通章節。 - Apache Arrow 跨語言互通(2025):
https://arrow.apache.org/,三個工具共同底層的官方文件。 - pandas 2.x 升級指南(2024):
https://pandas.pydata.org/docs/whatsnew/index.html,包含 PyArrow backend 的整合細節。 - DuckDB 查詢最佳化說明(1.4,2025):
https://duckdb.org/docs/,包含 predicate pushdown 與 projection pushdown 的設計。 - NYC 计程車示範資料(CC0 授權):
https://duckdb.org/data/nyc-taxi.csv.gz,本系列範例的來源。
留言
張貼留言