一. 前言:真正難的不是同時跑,而是怎麼收尾 #
假設你有 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:
InterpreterPoolExecutorExecutor.map(..., buffersize=...)
請注意:concurrent.futures.Future 和 asyncio.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],維持輸入順序
即使 3 比 1 先完成,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 有三條重要規則:
- worker 函式與參數必須能被
pickle。 - worker 必須能 import
__main__,不要期待 REPL 裡的 lambda 能正常送過去。 - 腳本入口要加
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 的 sys、builtins、__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。
比較安全的設計原則是:
- worker 做單一、封閉、可測試的函式。
- submit、等待、重試與取消都留在 coordinator。
- 不要讓 worker 偷偷建立巢狀 pool。
- 有相依性的流程,先在主執行緒建立 DAG 或分階段送出。
max_workers是資源上限,不是越大越快的魔法數字。
十一. 選擇清單:送出工作前先問七件事 #
- 工作主要在等 I/O,還是在算 CPU?
- callable、參數與回傳值能不能 pickle?
- task 是否依賴 mutable global state?
- 結果要維持輸入順序,還是完成就處理?
- 單筆失敗要保留其他結果,還是整批 fail fast?
- timeout 之後,task 如何合作式停止?
- 輸入量很大時,有沒有背壓與記憶體上限?
如果答案是「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。 並行不是免費加速;能清楚處理錯誤與收尾,才算真正跑得穩。🔭