Python concurrent.futures并發(fā)編程實戰(zhàn)

引言
在后端開發(fā)中,并發(fā)編程是提高系統(tǒng)性能的關(guān)鍵技術(shù)。Python的concurrent.futures模塊提供了簡潔的高級API,讓開發(fā)者能夠輕松實現(xiàn)多線程和多進程并發(fā)。作為一名從Python轉(zhuǎn)向Rust的后端開發(fā)者,我在實踐中總結(jié)了concurrent.futures的最佳實踐。本文將深入探討這個強大的并發(fā)工具,幫助你編寫高效的并發(fā)代碼。
一、concurrent.futures基礎(chǔ)
1.1 模塊概述
concurrent.futures模塊提供了兩個主要的執(zhí)行器:
- ThreadPoolExecutor:線程池執(zhí)行器
- ProcessPoolExecutor:進程池執(zhí)行器
1.2 基本使用模式
from concurrent.futures import ThreadPoolExecutor
def task(n):
return n * n
with ThreadPoolExecutor(max_workers=4) as executor:
future = executor.submit(task, 5)
result = future.result()
print(result)
1.3 核心組件
| 組件 | 說明 |
|---|---|
| Executor | 抽象執(zhí)行器基類 |
| ThreadPoolExecutor | 線程池執(zhí)行器 |
| ProcessPoolExecutor | 進程池執(zhí)行器 |
| Future | 表示異步計算的結(jié)果 |
二、ThreadPoolExecutor詳解
2.1 創(chuàng)建線程池
from concurrent.futures import ThreadPoolExecutor
# 創(chuàng)建線程池
executor = ThreadPoolExecutor(
max_workers=4,
thread_name_prefix='worker-'
)
# 使用上下文管理器
with ThreadPoolExecutor(max_workers=4) as executor:
# 提交任務(wù)
pass
2.2 提交任務(wù)
from concurrent.futures import ThreadPoolExecutor
def download_file(url):
# 模擬下載
import time
time.sleep(1)
return f"Downloaded: {url}"
with ThreadPoolExecutor(max_workers=3) as executor:
# 提交單個任務(wù)
future = executor.submit(download_file, "http://example.com/file1.txt")
# 獲取結(jié)果(阻塞)
result = future.result(timeout=5)
print(result)
2.3 批量提交任務(wù)
urls = [
"http://example.com/file1.txt",
"http://example.com/file2.txt",
"http://example.com/file3.txt",
]
with ThreadPoolExecutor(max_workers=3) as executor:
# 批量提交
futures = [executor.submit(download_file, url) for url in urls]
# 獲取所有結(jié)果
for future in futures:
print(future.result())
三、ProcessPoolExecutor詳解
3.1 創(chuàng)建進程池
from concurrent.futures import ProcessPoolExecutor
def compute_heavy(n):
# CPU密集型任務(wù)
return sum(i * i for i in range(n))
with ProcessPoolExecutor(max_workers=4) as executor:
future = executor.submit(compute_heavy, 1_000_000)
print(future.result())
3.2 線程池 vs 進程池
| 特性 | ThreadPoolExecutor | ProcessPoolExecutor |
|---|---|---|
| GIL限制 | 受GIL限制 | 不受GIL限制 |
| 適用場景 | IO密集型 | CPU密集型 |
| 啟動開銷 | 低 | 高 |
| 內(nèi)存開銷 | 低 | 高 |
| 數(shù)據(jù)共享 | 容易 | 困難 |
3.3 選擇建議
# IO密集型任務(wù) → ThreadPoolExecutor
# CPU密集型任務(wù) → ProcessPoolExecutor
# 混合任務(wù) → 結(jié)合使用
def process_task(data):
# IO操作
raw_data = fetch_from_api(data)
# CPU操作
result = compute(raw_data)
return result
四、Future對象詳解
4.1 Future狀態(tài)
from concurrent.futures import Future future = Future() # 檢查狀態(tài) print(future.done()) # False print(future.running()) # False print(future.cancelled()) # False # 設(shè)置結(jié)果 future.set_result(42) print(future.done()) # True print(future.result()) # 42
4.2 添加回調(diào)
def callback(future):
print(f"Task completed: {future.result()}")
with ThreadPoolExecutor() as executor:
future = executor.submit(task, 5)
future.add_done_callback(callback)
4.3 超時處理
try:
result = future.result(timeout=2)
except concurrent.futures.TimeoutError:
print("Task timed out")
五、高級用法
5.1 as_completed
from concurrent.futures import ThreadPoolExecutor, as_completed
def task(id):
import time
time.sleep(id)
return f"Task {id} completed"
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(task, i) for i in range(1, 4)]
# 按完成順序獲取結(jié)果
for future in as_completed(futures):
print(future.result())
5.2 map函數(shù)
with ThreadPoolExecutor(max_workers=3) as executor:
results = executor.map(task, [1, 2, 3, 4, 5])
# 按輸入順序返回結(jié)果
for result in results:
print(result)
5.3 wait函數(shù)
from concurrent.futures import wait, FIRST_COMPLETED
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(task, i) for i in range(1, 4)]
# 等待第一個完成
done, not_done = wait(futures, return_when=FIRST_COMPLETED)
print(f"Completed: {len(done)}")
print(f"Not completed: {len(not_done)}")
六、實戰(zhàn)案例
6.1 并行下載文件
import requests
from concurrent.futures import ThreadPoolExecutor
def download_file(url, save_path):
response = requests.get(url)
with open(save_path, 'wb') as f:
f.write(response.content)
return save_path
urls = [
("https://example.com/image1.jpg", "images/image1.jpg"),
("https://example.com/image2.jpg", "images/image2.jpg"),
("https://example.com/image3.jpg", "images/image3.jpg"),
]
with ThreadPoolExecutor(max_workers=5) as executor:
futures = [executor.submit(download_file, url, path) for url, path in urls]
for future in as_completed(futures):
print(f"Downloaded: {future.result()}")
6.2 并行數(shù)據(jù)庫查詢
import psycopg2
from concurrent.futures import ThreadPoolExecutor
def query_user(user_id):
conn = psycopg2.connect("dbname=example user=postgres")
cursor = conn.cursor()
cursor.execute("SELECT * FROM users WHERE id = %s", (user_id,))
result = cursor.fetchone()
conn.close()
return result
user_ids = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
with ThreadPoolExecutor(max_workers=4) as executor:
results = executor.map(query_user, user_ids)
for user_id, user in zip(user_ids, results):
print(f"User {user_id}: {user}")
6.3 混合IO和CPU任務(wù)
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def fetch_data(url):
import requests
return requests.get(url).json()
def process_data(data):
# CPU密集型處理
return sum(item['value'] for item in data)
def pipeline(url):
data = fetch_data(url)
return process_data(data)
urls = ["https://api.example.com/data1", "https://api.example.com/data2"]
# IO階段使用線程池
with ThreadPoolExecutor(max_workers=4) as io_executor:
futures = [io_executor.submit(fetch_data, url) for url in urls]
raw_data = [f.result() for f in futures]
# CPU階段使用進程池
with ProcessPoolExecutor(max_workers=4) as cpu_executor:
results = list(cpu_executor.map(process_data, raw_data))
print(results)
七、最佳實踐
7.1 合理設(shè)置worker數(shù)量
import os # CPU密集型任務(wù) cpu_workers = os.cpu_count() or 4 # IO密集型任務(wù) io_workers = min(32, (os.cpu_count() or 4) * 5)
7.2 避免共享狀態(tài)
# 不好的做法:共享可變狀態(tài)
counter = 0
def increment():
global counter
counter += 1
# 好的做法:使用線程安全的數(shù)據(jù)結(jié)構(gòu)或鎖
from threading import Lock
class ThreadSafeCounter:
def __init__(self):
self._count = 0
self._lock = Lock()
def increment(self):
with self._lock:
self._count += 1
7.3 優(yōu)雅關(guān)閉
executor = ThreadPoolExecutor(max_workers=4)
try:
# 提交任務(wù)
futures = [executor.submit(task, i) for i in range(10)]
# 獲取結(jié)果
for future in futures:
print(future.result())
finally:
# 關(guān)閉執(zhí)行器
executor.shutdown(wait=True)
八、性能對比
8.1 同步vs異步
import time
def sync_download(urls):
for url in urls:
download_file(url)
def async_download(urls):
with ThreadPoolExecutor(max_workers=5) as executor:
executor.map(download_file, urls)
# 性能對比
urls = ["https://example.com/file{}.txt".format(i) for i in range(10)]
start = time.time()
sync_download(urls)
print(f"Sync time: {time.time() - start:.2f}s")
start = time.time()
async_download(urls)
print(f"Async time: {time.time() - start:.2f}s")
8.2 線程池vs進程池
def cpu_intensive(n):
return sum(i * i for i in range(n))
# 線程池(受GIL限制)
with ThreadPoolExecutor(max_workers=4) as executor:
start = time.time()
executor.map(cpu_intensive, [10_000_000] * 4)
print(f"ThreadPool time: {time.time() - start:.2f}s")
# 進程池(不受GIL限制)
with ProcessPoolExecutor(max_workers=4) as executor:
start = time.time()
executor.map(cpu_intensive, [10_000_000] * 4)
print(f"ProcessPool time: {time.time() - start:.2f}s")
總結(jié)
concurrent.futures是Python并發(fā)編程的利器。通過本文的學(xué)習(xí),你應(yīng)該掌握了以下核心要點:
- ThreadPoolExecutor:適用于IO密集型任務(wù)
- ProcessPoolExecutor:適用于CPU密集型任務(wù)
- Future對象:管理異步計算結(jié)果
- 高級API:
as_completed、map、wait - 實戰(zhàn)案例:并行下載、數(shù)據(jù)庫查詢、混合任務(wù)
- 最佳實踐:合理設(shè)置worker數(shù)量、避免共享狀態(tài)
- 性能對比:同步vs異步、線程池vs進程池
作為從Python轉(zhuǎn)向Rust的后端開發(fā)者,理解并發(fā)編程模式對于構(gòu)建高性能系統(tǒng)至關(guān)重要。雖然Rust的并發(fā)模型更加安全,但Python的concurrent.futures提供了快速實現(xiàn)并發(fā)的便捷方式。
到此這篇關(guān)于Python concurrent.futures并發(fā)編程實戰(zhàn)的文章就介紹到這了,更多相關(guān)Python concurrent.futures并發(fā)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Python中JSON數(shù)據(jù)驗證的三種專業(yè)級方案
這篇文章主要介紹了Python中三種專業(yè)的JSON驗證方案:jsonschema、Pydantic和voluptuous,分別從聲明式驗證、類型安全的模型校驗和輕量級驗證流程三個方面進行了詳細(xì)講解,此外,還討論了傳統(tǒng)驗證方法的局限性,并展示了如何使用這些庫進行高效的數(shù)據(jù)校驗和模型定義2026-01-01
Python實現(xiàn)可獲取網(wǎng)易頁面所有文本信息的網(wǎng)易網(wǎng)絡(luò)爬蟲功能示例
這篇文章主要介紹了Python實現(xiàn)可獲取網(wǎng)易頁面所有文本信息的網(wǎng)易網(wǎng)絡(luò)爬蟲功能,涉及Python針對網(wǎng)頁的獲取、字符串正則判定等相關(guān)操作技巧,需要的朋友可以參考下2018-01-01
python實現(xiàn)自動登錄跳轉(zhuǎn)頁面并獲取信息
這篇文章主要為大家詳細(xì)介紹了如何使用python實現(xiàn)自動登錄跳轉(zhuǎn)頁面并獲取信息功能,文中的示例代碼講解詳細(xì),有需要的小伙伴可以跟隨小編一起學(xué)習(xí)一下2025-05-05
Python使用Pydantic驗證和解析配置數(shù)據(jù)的完整指南
在開發(fā)過程中,配置管理是繞不開的核心環(huán)節(jié),這些配置數(shù)據(jù)的質(zhì)量直接影響系統(tǒng)的穩(wěn)定性和安全性,下面我們就來看看如何使用Pydantic驗證和解析配置數(shù)據(jù)吧2026-01-01
python matplotlib餅狀圖參數(shù)及用法解析
這篇文章主要介紹了python matplotlib餅狀圖參數(shù)及用法解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下2019-11-11

