Python異步編程入門之實(shí)現(xiàn)文件批處理的并發(fā)處理方式
引言
在現(xiàn)代軟件開發(fā)中,處理大量文件或數(shù)據(jù)時(shí),提高處理效率和并發(fā)性是非常重要的。
Python 的 asyncio 庫提供了一種強(qiáng)大的方式來實(shí)現(xiàn)異步編程,從而提高程序的并發(fā)處理能力。
本文將面向 Python 初級(jí)程序員,介紹如何使用 asyncio 和 logging 模塊來實(shí)現(xiàn)一個(gè)異步批處理文件的并發(fā)處理系統(tǒng)。
代碼實(shí)現(xiàn)
1. 日志配置
首先,我們需要配置日志系統(tǒng),以便在處理文件時(shí)記錄日志信息。
日志配置包括設(shè)置日志格式和輸出位置。
import logging
import os
# 獲取當(dāng)前文件的絕對路徑
current_file = os.path.abspath(__file__)
# 配置日志格式
log_format = '%(asctime)s - %(levelname)s - %(pathname)s:%(lineno)d - %(message)s'
logging.basicConfig(format=log_format, level=logging.INFO)
# 創(chuàng)建一個(gè)文件處理器,并將日志輸出到文件
file_handler = logging.FileHandler('app.log')
file_handler.setFormatter(logging.Formatter(log_format))
logging.getLogger().addHandler(file_handler)2. 異步批處理類
接下來,我們定義一個(gè) AsyncBatchProcessor 類,用于處理批量文件。
該類使用 asyncio.Semaphore 來控制并發(fā)任務(wù)的數(shù)量。
import asyncio
import random
DEFAULT_MAX_CONCURRENT_TASKS = 2 # 最大并發(fā)任務(wù)數(shù)
MAX_RETRIES = 3 # 最大重試次數(shù)
class AsyncBatchProcessor:
def __init__(self, max_concurrent: int = DEFAULT_MAX_CONCURRENT_TASKS):
self.max_concurrent = max_concurrent
self.semaphore = asyncio.Semaphore(max_concurrent)
async def process_single_file(
self,
input_file: str,
retry_count: int = 0
) -> None:
"""處理單個(gè)文件的異步方法"""
async with self.semaphore: # 使用信號(hào)量控制并發(fā)
try:
logging.info(f"Processing file: {input_file}")
# 模擬文件處理過程
await asyncio.sleep(random.uniform(0.5, 2.0))
logging.info(f"Successfully processed {input_file}")
except Exception as e:
logging.error(f"Error processing {input_file} of Attempt {retry_count}: {str(e)}")
if retry_count < MAX_RETRIES:
logging.info(f"Retrying {input_file} (Attempt {retry_count + 1})")
await asyncio.sleep(1)
await self.process_single_file(input_file, retry_count + 1)
else:
logging.error(f"Failed to process {input_file} after {MAX_RETRIES} attempts")
async def process_batch(
self,
file_list: list
) -> None:
total_files = len(file_list)
logging.info(f"Found {total_files} files to process")
# 創(chuàng)建工作隊(duì)列
queue = asyncio.Queue()
# 將所有文件放入隊(duì)列
for file_path in file_list:
await queue.put(file_path)
# 創(chuàng)建工作協(xié)程
async def worker(worker_id: int):
while True:
try:
# 非阻塞方式獲取任務(wù)
input_file_path = await queue.get()
logging.info(f"Worker {worker_id} processing: {input_file_path}")
try:
await self.process_single_file(input_file_path)
except Exception as e:
logging.error(f"Error processing {input_file_path}: {str(e)}")
finally:
queue.task_done()
except asyncio.QueueEmpty:
# 隊(duì)列為空,工作結(jié)束
break
except Exception as e:
logging.error(f"Worker {worker_id} encountered error: {str(e)}")
break
# 創(chuàng)建工作任務(wù)
workers = []
for i in range(self.max_concurrent):
worker_task = asyncio.create_task(worker(i))
workers.append(worker_task)
# 等待隊(duì)列處理完成
await queue.join()
# 取消所有仍在運(yùn)行的工作任務(wù)
for w in workers:
w.cancel()
# 等待所有工作任務(wù)完成
await asyncio.gather(*workers, return_exceptions=True)3. 異步批處理入口函數(shù)
最后,我們定義一個(gè)異步批處理入口函數(shù) batch_detect,用于啟動(dòng)批處理任務(wù)。
async def batch_detect(
file_list: list,
max_concurrent: int = DEFAULT_MAX_CONCURRENT_TASKS
):
"""異步批處理入口函數(shù)"""
processor = AsyncBatchProcessor(max_concurrent)
await processor.process_batch(file_list)
# 示例調(diào)用
file_list = ["file1.pdf", "file2.pdf", "file3.pdf", "file4.pdf"]
asyncio.run(batch_detect(file_list))代碼解釋
1.日志配置:
- 使用
logging模塊記錄日志信息,包括時(shí)間、日志級(jí)別、文件路徑和行號(hào)、以及日志消息。 - 日志輸出到文件
app.log中,便于后續(xù)查看和分析。
2.異步批處理類 AsyncBatchProcessor:
__init__方法初始化最大并發(fā)任務(wù)數(shù)和信號(hào)量。process_single_file方法處理單個(gè)文件,使用信號(hào)量控制并發(fā),模擬文件處理過程,并在失敗時(shí)重試。process_batch方法處理批量文件,創(chuàng)建工作隊(duì)列和協(xié)程,控制并發(fā)任務(wù)的執(zhí)行。
3.異步批處理入口函數(shù) batch_detect:
- 創(chuàng)建
AsyncBatchProcessor實(shí)例,并調(diào)用process_batch方法啟動(dòng)批處理任務(wù)。
總結(jié)
通過使用 asyncio 和 logging 模塊,我們實(shí)現(xiàn)了一個(gè)高效的異步批處理文件系統(tǒng)。
該系統(tǒng)能夠并發(fā)處理大量文件,并在處理失敗時(shí)自動(dòng)重試,直到達(dá)到最大重試次數(shù)。
日志系統(tǒng)幫助我們記錄每個(gè)文件的處理過程,便于后續(xù)的調(diào)試和分析。
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
使用Python第三方庫pygame寫個(gè)貪吃蛇小游戲
這篇文章主要介紹了使用Python第三方庫pygame寫個(gè)貪吃蛇小游戲,本文通過實(shí)例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-03-03
python實(shí)現(xiàn)跨進(jìn)程(跨py文件)通信示例
本文主要介紹了python實(shí)現(xiàn)跨進(jìn)程(跨py文件)通信示例,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2022-03-03
jupyter notebook更換皮膚主題的實(shí)現(xiàn)
這篇文章主要介紹了jupyter notebook更換皮膚主題的實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2021-01-01
Pandas DataFrame進(jìn)行數(shù)據(jù)拼接方法詳解
這篇文章主要為大家詳細(xì)介紹了Pandas DataFrame進(jìn)行數(shù)據(jù)拼接多種方法,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2025-11-11
keras的load_model實(shí)現(xiàn)加載含有參數(shù)的自定義模型
這篇文章主要介紹了keras的load_model實(shí)現(xiàn)加載含有參數(shù)的自定義模型,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧2020-06-06
利用Python實(shí)現(xiàn)添加或讀取Excel公式
Excel公式是數(shù)據(jù)處理的核心工具,從簡單的加減運(yùn)算到復(fù)雜的邏輯判斷,掌握基礎(chǔ)語法是高效工作的起點(diǎn),下面我們就來看看如何使用Python進(jìn)行Excel公式的添加與讀取吧2025-03-03
深入淺出Python contextlib如何優(yōu)雅管理上下文資源
這篇文章主要為大家詳細(xì)介紹了Python contextlib優(yōu)雅管理上下文資源的相關(guān)知識(shí),文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2026-04-04
Python繪制分形圖案探索無限細(xì)節(jié)和奇妙之美
本文將介紹如何使用Python繪制各種分形圖案,包括分形樹、科赫曲線、曼德博集合等。通過本文讀者可以了解分形圖案的基本概念和構(gòu)造方法,并學(xué)會(huì)使用Python繪制出各種精美的分形圖案。本文還提供了具體的代碼示例和實(shí)踐案例,幫助讀者更好地理解分形圖案的奇妙之美2023-04-04

