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

引言
在后端開發(fā)中,并發(fā)編程是提高系統(tǒng)性能的關(guān)鍵技術(shù)。Python的concurrent.futures模塊提供了簡(jiǎn)潔的高級(jí)API,讓開發(fā)者能夠輕松實(shí)現(xiàn)多線程和多進(jìn)程并發(fā)。作為一名從Python轉(zhuǎn)向Rust的后端開發(fā)者,我在實(shí)踐中總結(jié)了concurrent.futures的最佳實(shí)踐。本文將深入探討這個(gè)強(qiáng)大的并發(fā)工具,幫助你編寫高效的并發(fā)代碼。
一、concurrent.futures基礎(chǔ)
1.1 模塊概述
concurrent.futures模塊提供了兩個(gè)主要的執(zhí)行器:
- ThreadPoolExecutor:線程池執(zhí)行器
- ProcessPoolExecutor:進(jìn)程池執(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 | 進(jìn)程池執(zhí)行器 |
| Future | 表示異步計(jì)算的結(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:
# 提交單個(gè)任務(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)建進(jìn)程池
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 進(jìn)程池
| 特性 | ThreadPoolExecutor | ProcessPoolExecutor |
|---|---|---|
| GIL限制 | 受GIL限制 | 不受GIL限制 |
| 適用場(chǎng)景 | IO密集型 | CPU密集型 |
| 啟動(dòng)開銷 | 低 | 高 |
| 內(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對(duì)象詳解
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 超時(shí)處理
try:
result = future.result(timeout=2)
except concurrent.futures.TimeoutError:
print("Task timed out")
五、高級(jí)用法
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)]
# 等待第一個(gè)完成
done, not_done = wait(futures, return_when=FIRST_COMPLETED)
print(f"Completed: {len(done)}")
print(f"Not completed: {len(not_done)}")
六、實(shí)戰(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階段使用進(jìn)程池
with ProcessPoolExecutor(max_workers=4) as cpu_executor:
results = list(cpu_executor.map(process_data, raw_data))
print(results)
七、最佳實(shí)踐
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)
八、性能對(duì)比
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)
# 性能對(duì)比
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進(jìn)程池
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")
# 進(jìn)程池(不受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)該掌握了以下核心要點(diǎn):
- ThreadPoolExecutor:適用于IO密集型任務(wù)
- ProcessPoolExecutor:適用于CPU密集型任務(wù)
- Future對(duì)象:管理異步計(jì)算結(jié)果
- 高級(jí)API:
as_completed、map、wait - 實(shí)戰(zhàn)案例:并行下載、數(shù)據(jù)庫查詢、混合任務(wù)
- 最佳實(shí)踐:合理設(shè)置worker數(shù)量、避免共享狀態(tài)
- 性能對(duì)比:同步vs異步、線程池vs進(jìn)程池
作為從Python轉(zhuǎn)向Rust的后端開發(fā)者,理解并發(fā)編程模式對(duì)于構(gòu)建高性能系統(tǒng)至關(guān)重要。雖然Rust的并發(fā)模型更加安全,但Python的concurrent.futures提供了快速實(shí)現(xiàn)并發(fā)的便捷方式。
到此這篇關(guān)于Python concurrent.futures并發(fā)編程實(shí)戰(zhàn)的文章就介紹到這了,更多相關(guān)Python concurrent.futures并發(fā)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
python自動(dòng)化測(cè)試通過日志3分鐘定位bug
軟件開發(fā)中通過日志記錄程序的運(yùn)行情況是一個(gè)開發(fā)的好習(xí)慣,對(duì)于錯(cuò)誤排查和系統(tǒng)運(yùn)維都有很大幫助,Python標(biāo)準(zhǔn)庫自帶了強(qiáng)大的logging日志模塊,在各種python模塊中得到廣泛應(yīng)用2021-11-11
python模塊簡(jiǎn)介之有序字典(OrderedDict)
字典是Python開發(fā)中很常用的一種數(shù)據(jù)結(jié)構(gòu),但dict有個(gè)缺陷(其實(shí)也不算缺陷),迭代時(shí)并不是按照元素添加的順序進(jìn)行,可能在某些場(chǎng)景下,不能滿足我們的要求。2016-12-12
Python中JSON數(shù)據(jù)驗(yàn)證的三種專業(yè)級(jí)方案
這篇文章主要介紹了Python中三種專業(yè)的JSON驗(yàn)證方案:jsonschema、Pydantic和voluptuous,分別從聲明式驗(yàn)證、類型安全的模型校驗(yàn)和輕量級(jí)驗(yàn)證流程三個(gè)方面進(jìn)行了詳細(xì)講解,此外,還討論了傳統(tǒng)驗(yàn)證方法的局限性,并展示了如何使用這些庫進(jìn)行高效的數(shù)據(jù)校驗(yàn)和模型定義2026-01-01
Python實(shí)現(xiàn)可獲取網(wǎng)易頁面所有文本信息的網(wǎng)易網(wǎng)絡(luò)爬蟲功能示例
這篇文章主要介紹了Python實(shí)現(xiàn)可獲取網(wǎng)易頁面所有文本信息的網(wǎng)易網(wǎng)絡(luò)爬蟲功能,涉及Python針對(duì)網(wǎng)頁的獲取、字符串正則判定等相關(guān)操作技巧,需要的朋友可以參考下2018-01-01
python實(shí)現(xiàn)自動(dòng)登錄跳轉(zhuǎn)頁面并獲取信息
這篇文章主要為大家詳細(xì)介紹了如何使用python實(shí)現(xiàn)自動(dòng)登錄跳轉(zhuǎn)頁面并獲取信息功能,文中的示例代碼講解詳細(xì),有需要的小伙伴可以跟隨小編一起學(xué)習(xí)一下2025-05-05
Python使用Pydantic驗(yàn)證和解析配置數(shù)據(jù)的完整指南
在開發(fā)過程中,配置管理是繞不開的核心環(huán)節(jié),這些配置數(shù)據(jù)的質(zhì)量直接影響系統(tǒng)的穩(wěn)定性和安全性,下面我們就來看看如何使用Pydantic驗(yàn)證和解析配置數(shù)據(jù)吧2026-01-01
python matplotlib餅狀圖參數(shù)及用法解析
這篇文章主要介紹了python matplotlib餅狀圖參數(shù)及用法解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-11-11

