淺析Python 3.11以下如何優(yōu)雅地實(shí)現(xiàn)自動取消任務(wù)
?小李今天又遇到了煩心事。
他寫了一個(gè)數(shù)據(jù)處理腳本,要調(diào)用外部API獲取一萬個(gè)用戶的信息。每個(gè)請求大概要等2秒。他不想干等著,所以用了asyncio,并發(fā)發(fā)出去50個(gè)請求。
跑了三分鐘,腳本卡住了。
不是死鎖,就是單純的慢——有幾個(gè)API服務(wù)不穩(wěn)定,響應(yīng)要半分鐘。他想:“能不能設(shè)定一個(gè)超時(shí)時(shí)間,比如5秒,超過5秒就不等了,直接跳過?”
這個(gè)需求太常見了。不管是網(wǎng)絡(luò)請求、數(shù)據(jù)庫查詢還是復(fù)雜的計(jì)算任務(wù),你總不希望它無限期地卡下去。就像點(diǎn)外賣,等了45分鐘還沒到,你肯定想取消訂單,換一家點(diǎn)。
問題在于,Python的異步超時(shí)機(jī)制,在3.11之前有一個(gè)不大不小的“坑”,如果你不注意,任務(wù)可能根本取消不掉。
一個(gè)“取消失敗”的真實(shí)案例
先來看一段代碼,它在Python 3.10上運(yùn)行:
import asyncio
async def quick_task():
# 這個(gè)任務(wù)瞬間完成
return "done"
async def wrapper():
# 設(shè)置30秒超時(shí),但實(shí)際任務(wù)1秒就完成
return await asyncio.wait_for(quick_task(), timeout=30)
async def main():
task = asyncio.create_task(wrapper())
await asyncio.sleep(0) # 讓任務(wù)開始執(zhí)行
task.cancel() # 手動取消
try:
await task
except asyncio.CancelledError:
print("任務(wù)被取消了")
else:
print("任務(wù)沒有被取消!")
asyncio.run(main())
在Python 3.8上運(yùn)行,輸出的是“任務(wù)被取消了”。但在Python 3.9或3.10上運(yùn)行,輸出的卻是“任務(wù)沒有被取消!”
什么情況?明明調(diào)用了cancel(),為什么任務(wù)沒被取消?
真相藏在wait_for的源碼里
要理解這個(gè)問題,得看看asyncio.wait_for這個(gè)函數(shù)到底在干什么。
wait_for的作用是:給一個(gè)任務(wù)設(shè)置一個(gè)超時(shí)時(shí)間,如果超時(shí)了就取消它。但這里有一個(gè)細(xì)節(jié)——如果任務(wù)在超時(shí)之前就已經(jīng)完成了,而在完成的那一瞬間你恰好發(fā)起了取消請求,wait_for會怎么處理?
Python 3.8的做法是:不管任務(wù)完沒完成,只要收到了取消信號,就拋出一個(gè)CancelledError。
Python 3.9及之后的做法是:先檢查一下任務(wù)是不是已經(jīng)完成了。如果已經(jīng)完成了,“取消”就沒有意義了,直接返回任務(wù)的結(jié)果,不拋異常。
聽起來3.9之后的邏輯更合理對吧?畢竟任務(wù)都做完了,還取消什么?
但在并發(fā)環(huán)境下,問題就出在這個(gè)“先檢查”上。
想象一下這個(gè)時(shí)序:
- 任務(wù)在
wait_for里等著 - 你調(diào)用
task.cancel() wait_for收到取消信號- 它問:任務(wù)完成了嗎?如果回答“完成了”,它就直接返回結(jié)果,把取消信號吞掉
這個(gè)“任務(wù)完成了嗎”的判斷,在多線程/多任務(wù)的并發(fā)環(huán)境下,會因?yàn)闀r(shí)序問題出現(xiàn)誤判。任務(wù)實(shí)際上還沒有真正完成,只是處于“即將完成”的狀態(tài),wait_for就可能把它當(dāng)成“已完成”,從而無視你的取消請求。
這不是bug,而是設(shè)計(jì)上的一種取舍。但這個(gè)取舍導(dǎo)致了一個(gè)后果:在Python 3.11以下的版本中,取消操作并不是100%可靠的。
那怎么辦?三個(gè)靠譜的解決方案
既然wait_for靠不住,那就自己動手。
方案一:自己封裝一個(gè)“靠譜版的wait_for”
思路很簡單:不依賴wait_for的取消機(jī)制,而是自己用asyncio.create_task()加一個(gè)超時(shí)的“看門狗”。
import asyncio
async def cancellable_wait_for(coro, timeout):
"""
一個(gè)更可靠的wait_for版本,確保取消信號不會被吞掉
"""
task = asyncio.create_task(coro)
try:
# 等待任務(wù)完成,或者超時(shí)
return await asyncio.wait_for(task, timeout=timeout)
except asyncio.TimeoutError:
# 超時(shí)了,取消任務(wù)
task.cancel()
try:
await task
except asyncio.CancelledError:
# 確保取消信號被傳播出去
pass
raise # 重新拋出TimeoutError
except asyncio.CancelledError:
# 外部取消了,把取消信號傳遞給內(nèi)部任務(wù)
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
raise # 重新拋出CancelledError
async def my_task():
try:
await asyncio.sleep(10)
return "完成"
except asyncio.CancelledError:
print("內(nèi)部任務(wù)被取消了")
raise
async def main():
task = asyncio.create_task(cancellable_wait_for(my_task(), timeout=5))
await asyncio.sleep(2)
task.cancel()
try:
await task
except asyncio.CancelledError:
print("外部:確實(shí)被取消了")
asyncio.run(main())
這個(gè)方案的要點(diǎn)是:不管什么情況,只要外部取消或者超時(shí),都強(qiáng)制取消內(nèi)部任務(wù),然后把取消信號往上拋。不會被“任務(wù)已完成”這種假象迷惑。
方案二:用第三方庫quattro的CancelScope
Python 3.11之后,官方推出了asyncio.timeout(),終于有了靠譜的超時(shí)機(jī)制。那3.11以下怎么辦?
有一個(gè)第三方庫叫quattro,它在3.11以下版本中實(shí)現(xiàn)了類似的功能。
# pip install quattro
import asyncio
from quattro import move_on_after
async def main():
with move_on_after(5) as scope: # 5秒超時(shí)
result = await some_slow_operation()
print(result)
if scope.cancelled_caught:
print("超時(shí)被取消了,但程序繼續(xù)運(yùn)行")
asyncio.run(main())
quattro的CancelScope比官方的更靈活:它是普通的上下文管理器(不需要async with),可以手動調(diào)用scope.cancel()提前取消,還可以查詢是否被取消了。
如果你在項(xiàng)目中需要同時(shí)支持Python 3.9、3.10和3.11,quattro是個(gè)不錯的選擇。代碼寫一遍,在所有版本上行為一致。
方案三:最底層的做法——自己用Event和超時(shí)循環(huán)
如果不想引入第三方庫,也嫌自己封裝太麻煩,還有一個(gè)最樸素的辦法:不用wait_for,自己用asyncio.Event加上循環(huán)檢查。
import asyncio
async def cancellable_operation(timeout):
"""
在操作內(nèi)部主動檢查超時(shí)和取消信號
"""
start = asyncio.get_event_loop().time()
# 假設(shè)這是一系列的小步驟
for step in range(10):
# 檢查是否超時(shí)
if asyncio.get_event_loop().time() - start > timeout:
raise asyncio.TimeoutError()
# 檢查是否被取消(通過捕獲取消信號)
try:
await asyncio.sleep(0.5) # 模擬一個(gè)小步驟
except asyncio.CancelledError:
# 做清理工作
print("收到取消信號,正在清理...")
raise # 重新拋出,讓上層知道被取消了
print(f"完成步驟 {step}")
return "全部完成"
async def main():
task = asyncio.create_task(cancellable_operation(timeout=3))
await asyncio.sleep(2)
task.cancel()
try:
result = await task
print(result)
except asyncio.TimeoutError:
print("超時(shí)了")
except asyncio.CancelledError:
print("被取消了")
asyncio.run(main())
這個(gè)方案的優(yōu)點(diǎn)是:你完全掌控取消邏輯。缺點(diǎn)是:你得在操作內(nèi)部主動插入檢查點(diǎn)。如果你的操作是一大塊無法分割的同步代碼,這個(gè)方法就不太適用了。
那同步代碼怎么辦?線程池里怎么取消?
上面聊的都是異步代碼(async/await)。但小李的腳本里,requests.get()是同步的,它不響應(yīng)asyncio的取消信號。
這種情況下,task.cancel()根本沒用。因?yàn)槿∠盘栔辉?code>await的地方才能被處理,而同步代碼里沒有await。
解決方案是:用loop.run_in_executor把同步代碼扔到線程池里,然后用一個(gè)threading.Event來手動控制停止。
import asyncio
import threading
import time
def blocking_task(stop_event):
"""
這是一個(gè)會阻塞的同步函數(shù)
它定期檢查stop_event,如果被設(shè)置了就主動退出
"""
print("同步任務(wù)開始")
for i in range(30):
if stop_event.is_set():
print("收到停止信號,主動退出")
return "被停止了"
print(f"工作中... {i}")
time.sleep(1) # 模擬耗時(shí)操作
return "正常完成"
async def main():
stop_event = threading.Event()
loop = asyncio.get_running_loop()
# 把同步任務(wù)扔到線程池里執(zhí)行
task = loop.run_in_executor(None, blocking_task, stop_event)
# 5秒后發(fā)送停止信號
await asyncio.sleep(5)
stop_event.set()
# 等待任務(wù)結(jié)束
result = await task
print(f"結(jié)果: {result}")
asyncio.run(main())
關(guān)鍵點(diǎn):同步任務(wù)本身必須主動檢查stop_event,不能指望Python幫你“強(qiáng)行”取消。強(qiáng)行殺線程在Python里是不安全的,也不被推薦。
那如果用多進(jìn)程呢?
如果你面對的是CPU密集型的任務(wù),線程也幫不了你——Python的GIL鎖會讓多線程在計(jì)算任務(wù)上毫無優(yōu)勢。這時(shí)候考慮用ProcessPoolExecutor:
import asyncio
from concurrent.futures import ProcessPoolExecutor
import time
def cpu_intensive_task():
"""模擬一個(gè)費(fèi)CPU的大活兒"""
total = 0
for i in range(100000000):
total += i
# 每1000萬次檢查一次,但進(jìn)程間通信復(fù)雜,先忽略細(xì)節(jié)
if i % 10000000 == 0:
print(f"計(jì)算到 {i}")
return total
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor(max_workers=1) as pool:
# 在子進(jìn)程中執(zhí)行
task = loop.run_in_executor(pool, cpu_intensive_task)
# 等待3秒
await asyncio.sleep(3)
print("3秒到了,但子進(jìn)程不會自己停...")
# 注意:ProcessPoolExecutor沒有優(yōu)雅的取消方法
# 只能shutdown(wait=False)但會有warning
# 或者terminate,但不推薦
多進(jìn)程的取消更麻煩。進(jìn)程不像線程,你不能優(yōu)雅地通知它“停下來”。常見的做法是:在子進(jìn)程里也放一個(gè)檢查循環(huán),通過進(jìn)程間通信(比如multiprocessing.Event或隊(duì)列)來傳遞停止信號。
總結(jié)一下,到底選哪個(gè)方案?
給你一張決策表:
| 你的場景 | 推薦方案 |
|---|---|
| 純異步代碼,Python 3.11+ | 直接用官方asyncio.timeout() |
| 純異步代碼,Python 3.10及以下 | 用quattro的CancelScope,或自己封裝cancellable_wait_for |
| 同步阻塞代碼(如requests) | 線程池 + threading.Event手動輪詢檢查 |
| CPU密集型代碼 | 多進(jìn)程 + 進(jìn)程間通信信號,或用concurrent.futures的超時(shí)(有限支持) |
| 不想改現(xiàn)有代碼,只想加個(gè)保護(hù) | 用asyncio.wait_for,但要接受它可能偶爾失效 |
小李最后選擇了方案一:自己封裝了一個(gè)cancellable_wait_for,用了半小時(shí)寫完,以后所有異步任務(wù)都走這個(gè)函數(shù),再也沒出現(xiàn)過“取消不掉”的情況。
他說了一句大實(shí)話:“官方的不靠譜,就自己寫一個(gè)。反正常用的功能就那么幾個(gè),封裝一次到處用,不虧。”
到此這篇關(guān)于淺析Python 3.11以下如何優(yōu)雅地實(shí)現(xiàn)自動取消任務(wù)的文章就介紹到這了,更多相關(guān)Python自動取消任務(wù)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Python爬蟲實(shí)現(xiàn)爬取下載網(wǎng)站數(shù)據(jù)的幾種方法示例
這篇文章主要為大家介紹了Python爬蟲實(shí)現(xiàn)爬取下載網(wǎng)站數(shù)據(jù)的幾種方法示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-11-11
Python Web服務(wù)器Tornado使用小結(jié)
最近在做一個(gè)網(wǎng)站的后端開發(fā)。因?yàn)槌跗谥挥形乙粋€(gè)人做,所以技術(shù)選擇上很自由。在 web 服務(wù)器上我選擇了 Tornado。雖然曾經(jīng)也讀過它的源碼,并做過一些小的 demo,但畢竟這是第一次在工作中使用,難免又發(fā)現(xiàn)了一些值得分享的東西2014-05-05
解決python xx.py文件點(diǎn)擊完之后一閃而過的問題
今天小編就為大家分享一篇解決python xx.py文件點(diǎn)擊完之后一閃而過的問題,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧2019-06-06
python腳本設(shè)置系統(tǒng)時(shí)間的兩種方法
這篇文章主要介紹了python腳本設(shè)置系統(tǒng)時(shí)間的兩種方法,其一是調(diào)用socket直接發(fā)送udp包到國家授時(shí)中心,其二是調(diào)用ntplib包,感興趣的小伙伴們可以參考一下2016-02-02
opencv檢測動態(tài)物體的實(shí)現(xiàn)
本文主要介紹了opencv檢測動態(tài)物體的實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2021-07-07
Python使用Pyecharts繪制交互式中國地圖的流程解析
Pyecharts作為基于ECharts的Python可視化庫,提供三大核心優(yōu)勢,即零代碼交互,多級地圖支持和高度定制化,下面我們就來看看如何使用Pyecharts繪制交互式中國地圖吧2026-02-02
pandas按行按列遍歷Dataframe的三種方式小結(jié)
本文主要介紹了pandas按行按列遍歷Dataframe,主要介紹了三種方法,具有一定的參考價(jià)值,感興趣的可以了解一下2023-11-11
如何使用python實(shí)現(xiàn)多個(gè)csv文件數(shù)據(jù)的合并和輸出
文章介紹了如何使用Python批量合并多個(gè)CSV文件,并提供具體代碼示例,代碼簡單易懂,感興趣的朋友一起看看吧2025-03-03
python實(shí)現(xiàn)DNS正向查詢、反向查詢的例子
這篇文章主要介紹了python實(shí)現(xiàn)DNS正向查詢、反向查詢的例子,需要的朋友可以參考下2014-04-04

