快轉到主要內容
  1. 教學文章/

Python concurrent.futures 實戰:Thread、Process 與 Interpreter Pool

·8 分鐘· loading · loading · ·
Python Concurrent.futures ThreadPoolExecutor ProcessPoolExecutor InterpreterPoolExecutor Concurrency Parallelism
每日拍拍
作者
每日拍拍
科學家 X 科技宅宅
目錄
Python 學習 - 本文屬於一個選集。
§ 122: 本文

一. 前言:真正難的不是同時跑,而是怎麼收尾
#

假設你有 500 個檔案要讀、50 張圖片要轉換,或一批 API 要呼叫。 把工作「同時丟出去」不難;麻煩的是後面這一串問題:

  • 哪個任務先完成?
  • 某個任務失敗時,其他結果還要不要收?
  • 等太久可以取消嗎?
  • I/O 與 CPU 工作該用同一種 pool 嗎?
  • 程式結束時,worker 有沒有真的清乾淨?

標準庫 concurrent.futures 把這些問題整理成兩個核心抽象:

  • Executor:負責排程與執行工作。
  • Future:代表一個現在還不知道、稍後才會出現的結果。

同一套介面可以接上 thread、process,Python 3.14 還加入了 interpreter pool。 拍拍君今天不從鎖、Queue 或 coroutine 重新教起, 而是專注在「如何把一批工作可靠地送出、觀察、停止與回收」。

如果你想先補底層觀念,可以看:

二. 安裝與版本:它就在標準庫裡
#

concurrent.futures 從 Python 3.2 起就是標準庫,不需要 pip install

python --version
python -c "import concurrent.futures; print('ready')"

本文一般範例可在近年的 Python 3 執行;以下兩項需要 Python 3.14:

  • InterpreterPoolExecutor
  • Executor.map(..., buffersize=...)

請注意:concurrent.futures.Futureasyncio.Future 不是同一個類別。 前者由 Executor 產生,後者屬於 event loop 的 coroutine 世界。

三. 核心模型:submit 之後拿到 Future
#

先用最小範例認識介面:

from concurrent.futures import ThreadPoolExecutor
import time
def normalize_name(name: str) -> str:
    time.sleep(0.05)
    return name.strip().title()
with ThreadPoolExecutor(max_workers=2) as executor:
    future = executor.submit(normalize_name, "  拍拍醬  ")
    print(future.running() or future.done())
    print(future.result(timeout=1))

submit(fn, *args, **kwargs) 會立即回傳 Future, 工作可能還在 queue、正在執行,或已經完成。

最常用的狀態方法是:

方法 意義
running() 已開始執行,通常不能取消
done() 已完成、失敗,或成功取消
cancel() 嘗試取消尚未開始的工作
cancelled() 是否真的取消成功
result(timeout=...) 取得結果;失敗時重新拋出原例外
exception(timeout=...) 取得例外物件;成功時回傳 None

Future 不是背景工作的「遙控器」。 cancel() 只能阻止尚未開始的任務,不能安全地把正在跑的 Python 函式切成兩半。

四. map 還是 submit:先決定你要哪種結果順序
#

當輸入與輸出是一對一時,map() 最簡潔:

from concurrent.futures import ThreadPoolExecutor
import time
def delayed_square(number: int) -> int:
    time.sleep((4 - number) * 0.02)
    return number * number
with ThreadPoolExecutor(max_workers=3) as executor:
    results = list(executor.map(delayed_square, [1, 2, 3]))
print(results)  # [1, 4, 9],維持輸入順序

即使 31 先完成,map() 仍按輸入順序 yield 結果。 這對報表很方便,但慢任務也可能擋住後面已完成的結果。

如果你要「誰先完成就先處理誰」,改用 submit()as_completed()

from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def fetch(item: str, delay: float) -> str:
    time.sleep(delay)
    return f"{item}:ok"
jobs = {"alpha": 0.08, "beta": 0.02, "gamma": 0.05}
with ThreadPoolExecutor(max_workers=3) as executor:
    future_to_name = {
        executor.submit(fetch, name, delay): name
        for name, delay in jobs.items()
    }
    for future in as_completed(future_to_name):
        name = future_to_name[future]
        print(name, future.result())

拍拍君的簡化選擇規則:

  • 只要順序一致、每筆處理方式相同:用 map()
  • 要逐筆處理例外、顯示進度或保存 metadata:用 submit()
  • 想先顯示最快完成的結果:搭配 as_completed()

五. 三種 Executor:先看工作,不要先看名字
#

Python 3.14 的三種 pool 共用 Executor / Future 介面, 但隔離、啟動成本與資料傳遞方式完全不同。

Executor 適合 真正多核心跑 Python bytecode 主要代價
ThreadPoolExecutor 網路、磁碟、等待型 I/O 通常否 共享狀態與競爭條件
ProcessPoolExecutor CPU 密集、隔離需求高 啟動、序列化、額外記憶體
InterpreterPoolExecutor 可 pickle、可隔離的 CPU 工作 Python 3.14、模組與狀態隔離

5.1 I/O 工作:ThreadPoolExecutor
#

Thread 共享同一個 process 與記憶體,傳參數很方便。 當工作主要等待 HTTP、檔案或資料庫時,它通常是最簡單的起點。 代價是可變共享狀態需要同步,而且純 Python CPU 工作通常受 GIL 限制。

5.2 CPU 工作:ProcessPoolExecutor
#

CPU 密集工作可以交給獨立 process:

from concurrent.futures import ProcessPoolExecutor
import math
def is_prime(number: int) -> bool:
    if number < 2:
        return False
    if number == 2:
        return True
    if number % 2 == 0:
        return False
    limit = math.isqrt(number)
    return all(number % divisor for divisor in range(3, limit + 1, 2))
def main() -> None:
    numbers = [999_983, 1_000_003, 1_000_033]
    with ProcessPoolExecutor() as executor:
        results = executor.map(is_prime, numbers, chunksize=1)
        print(list(zip(numbers, results)))
if __name__ == "__main__":
    main()

Process pool 有三條重要規則:

  1. worker 函式與參數必須能被 pickle
  2. worker 必須能 import __main__,不要期待 REPL 裡的 lambda 能正常送過去。
  3. 腳本入口要加 if __name__ == "__main__":

Python 3.14 的預設 process start method 已不再是 fork。 不要依賴「子程序剛好繼承父程序所有狀態」;明確傳入資料會更可攜。

5.3 Python 3.14:InterpreterPoolExecutor
#

Interpreter pool 的每個 worker 都是一條 thread, 但每條 thread 擁有獨立 interpreter 與自己的 GIL,因此能使用多核心。

import concurrent.futures as futures
import math
def count_factors(number: int) -> int:
    limit = math.isqrt(number)
    count = 0
    for candidate in range(1, limit + 1):
        if number % candidate == 0:
            count += 1 if candidate * candidate == number else 2
    return count
def main() -> None:
    executor_type = getattr(futures, "InterpreterPoolExecutor", None)
    if executor_type is None:
        raise RuntimeError("這個範例需要 Python 3.14+")
    numbers = [10_000_019, 10_000_079, 10_000_103]
    with executor_type(max_workers=3) as executor:
        print(list(executor.map(count_factors, numbers)))
if __name__ == "__main__":
    main()

它不是「不用 process 就能透明共享所有 Python 物件」。 每個 interpreter 的 sysbuiltins__main__ 與 imported module 都各自獨立; 送入 callable、參數與傳回值時仍會使用 pickle。

這表示它適合邊界清楚的工作,不適合依賴大量隱含 global state 的函式。 第三方 C extension 是否支援多 interpreter,也要在真實環境測試。

六. 例外不要消失:在 result() 的邊界收集
#

worker 裡的例外不會在 submit() 當下出現, 而是在呼叫 result()、走訪 map() 結果時重新拋出。

from concurrent.futures import ThreadPoolExecutor, as_completed
def parse_score(raw: str) -> int:
    score = int(raw)
    if not 0 <= score <= 100:
        raise ValueError(f"score out of range: {score}")
    return score
records = ["80", "oops", "105", "92"]
successes: list[tuple[str, int]] = []
failures: list[tuple[str, str]] = []
with ThreadPoolExecutor(max_workers=3) as executor:
    future_to_raw = {
        executor.submit(parse_score, raw): raw for raw in records
    }
    for future in as_completed(future_to_raw):
        raw = future_to_raw[future]
        try:
            successes.append((raw, future.result()))
        except (TypeError, ValueError) as exc:
            failures.append((raw, str(exc)))
print("ok", sorted(successes))
print("failed", sorted(failures))

不要只寫一個寬廣的 except Exception: pass。 至少保留輸入識別、例外型別與訊息,否則批次工作看似成功,資料卻偷偷少了一半。

七. Timeout 只是停止等待,不是停止工作
#

這是 Future 最常被誤解的地方。

from concurrent.futures import ThreadPoolExecutor, TimeoutError
import time
def slow_job() -> str:
    time.sleep(0.2)
    return "finished"
with ThreadPoolExecutor(max_workers=1) as executor:
    future = executor.submit(slow_job)
    try:
        print(future.result(timeout=0.05))
    except TimeoutError:
        print("等待逾時,但工作可能仍在背景執行")
    print(future.result(timeout=1))

result(timeout=...) 控制的是呼叫端願意等多久, 不會中斷已經開始的 slow_job()

真正需要合作式取消時,讓 task 定期檢查停止訊號並主動 return。 threading.Event 適合 thread;process 或 interpreter 則需要各自相容的通訊方式。

八. 一個失敗就停:wait(FIRST_EXCEPTION)
#

有些批次可以保留部分成功,有些則必須 fail fast。 wait() 可以把這個政策寫得很明確:

from concurrent.futures import FIRST_EXCEPTION, ThreadPoolExecutor, wait
import time
def transform(value: int) -> int:
    time.sleep(0.01)
    if value == 3:
        raise ValueError("bad input: 3")
    return value * 10
with ThreadPoolExecutor(max_workers=2) as executor:
    futures = [executor.submit(transform, value) for value in range(8)]
    done, not_done = wait(futures, return_when=FIRST_EXCEPTION)
    failed = [future for future in done if future.exception() is not None]
    if failed:
        for future in not_done:
            future.cancel()  # 只能取消尚未開始者
        print("first error:", failed[0].exception())

離開 with 時,Executor 預設仍會等待已開始的任務結束。 若你自己管理生命週期,可用:

executor.shutdown(wait=True, cancel_futures=True)

這會取消還在排隊的 Future,但已開始的工作仍會完成。

Python 3.14 的 ProcessPoolExecutor 另外提供 terminate_workers()kill_workers(); 它們是破壞性較高的緊急手段,不該取代正常的 timeout、合作式停止與 finally 清理。

九. 控制背壓:Python 3.14 的 buffersize
#

舊版 Executor.map() 會很快收集 iterable 並送出大量任務。 輸入是百萬筆資料時,queue 與 Future 本身就可能吃掉可觀記憶體。

Python 3.14 可以設定 buffersize

from concurrent.futures import ThreadPoolExecutor
def normalize(number: int) -> str:
    return f"item-{number:06d}"
def huge_stream():
    for number in range(1_000_000):
        yield number
with ThreadPoolExecutor(max_workers=8) as executor:
    results = executor.map(normalize, huge_stream(), buffersize=32)
    for result in results:
        if result.endswith("000099"):
            print(result)
            break

當尚未 yield 的結果數量碰到 buffer 上限,輸入迭代會暫停。 這就是背壓:消費端跟不上時,不再無限制地塞工作。

另外,chunksize 只會影響 ProcessPoolExecutor; 對 Thread 與 Interpreter pool 沒有效果。

十. 避免兩種經典 Deadlock
#

第一種:pool 裡的 task 等同一個 pool 的另一個 Future, 但所有 worker 都已經被等待者占滿。

# 危險概念:max_workers=1 時,外層工作占住唯一 worker,
# 又 submit 內層工作並呼叫 inner.result(),內層永遠沒機會開始。

第二種:ProcessPoolExecutor 的 worker 內再次呼叫同一個 Executor 或 Future 方法。 官方文件直接警告這會造成 deadlock。

比較安全的設計原則是:

  1. worker 做單一、封閉、可測試的函式。
  2. submit、等待、重試與取消都留在 coordinator。
  3. 不要讓 worker 偷偷建立巢狀 pool。
  4. 有相依性的流程,先在主執行緒建立 DAG 或分階段送出。
  5. max_workers 是資源上限,不是越大越快的魔法數字。

十一. 選擇清單:送出工作前先問七件事
#

  1. 工作主要在等 I/O,還是在算 CPU?
  2. callable、參數與回傳值能不能 pickle?
  3. task 是否依賴 mutable global state?
  4. 結果要維持輸入順序,還是完成就處理?
  5. 單筆失敗要保留其他結果,還是整批 fail fast?
  6. timeout 之後,task 如何合作式停止?
  7. 輸入量很大時,有沒有背壓與記憶體上限?

如果答案是「I/O、多筆獨立、要逐筆結果」,先從 ThreadPool 開始。 如果是「純 Python CPU、多筆獨立、可 pickle」,比較 Process 與 Interpreter pool。 如果工作本來就是大量 socket coroutine,通常該回到 asyncio / AnyIO, 而不是為了使用 Future 硬塞進 thread。

結語
#

concurrent.futures 的價值不是幫你把 for 迴圈換成平行版而已, 而是提供一套共同語言描述任務的生命週期。

記住四個重點:

  • Executor 決定工作在哪裡跑,Future 描述結果的狀態。
  • map() 保持輸入順序,as_completed() 提供完成順序。
  • timeout 不會自動殺掉已開始的 task,取消政策要自己設計。
  • Thread、Process、Interpreter pool 介面相似,隔離與成本卻不同。

先用最小資料量量測,再調整 worker 數、chunksize 或 buffersize。 並行不是免費加速;能清楚處理錯誤與收尾,才算真正跑得穩。🔭

延伸閱讀
#

Python 學習 - 本文屬於一個選集。
§ 122: 本文

相關文章

DuckDB Window Functions 實戰:排名、移動平均與分組分析
·8 分鐘· loading · loading
Python DuckDB SQL Window Functions Data-Analysis Analytics OLAP
Python Requests 實戰:Session、Timeout、Retry 與可靠 API Client
·7 分鐘· loading · loading
Python Requests HTTP API Session Retry Developer-Tools
Pandas 資料清理實戰:缺失值、型別、重複資料與驗證
·7 分鐘· loading · loading
Python Pandas Data-Cleaning Missing-Data Validation ETL
scikit-learn Pipeline 實戰:前處理、交叉驗證與避免資料洩漏
·7 分鐘· loading · loading
Python Scikit-Learn Pipeline ColumnTransformer Machine-Learning Cross-Validation Data-Leakage
NumPy Broadcasting 實戰:Shape、索引、Mask 與維度對齊
·6 分鐘· loading · loading
Python Numpy Broadcasting Indexing Boolean Mask Data-Analysis
uv pip-compile migration:從 requirements.txt 到可重現部署流程
·7 分鐘· loading · loading
Python Uv Requirements.txt Pip-Compile Dependency-Management Deployment