در این فصل، بهصورت کامل و مستقل به بررسی مدیریت خطا در محیط multiprocessing میپردازیم. فرض میکنیم که شما با مفاهیم پایهی multiprocessing و فرایندها (Processes) آشنایی دارید و قصد دارید برنامههایی مقاوم و قابلاعتماد با استفاده از پردازشهای موازی بنویسید.
چالشهای منحصربهفرد multiprocessing در مدیریت خطا
multiprocessing در پایتون، با وجود قدرت بالا در استفاده از چندین هستهی پردازشی، چالشهای خاص خود را در مدیریت خطا دارد. بیایید ابتدا این چالشها را بشناسیم:
۱. حافظهی جداگانه
برخلاف threading که تردها حافظهی مشترک دارند، در multiprocessing هر پردازش حافظهی خودش را دارد. این یعنی خطاهای مربوط به حافظهی مشترک کمتر است، اما ارتباط بین پردازشها پیچیدهتر میشود.
۲. سربار بالای ایجاد پردازش
ایجاد یک پردازش جدید، سنگینتر از ایجاد یک ترد است. بنابراین در طراحی Retry، باید تعداد تلاشها را کمتر در نظر بگیرید و زمانهای انتظار را طولانیتر تنظیم کنید.
۳. مدیریت استثناها در پردازشها
اگر یک استثنا در یک پردازش رخ دهد و مدیریت نشود، آن پردازش از کار میافتد اما سایر پردازشها به کار خود ادامه میدهند. انتقال استثناها از پردازشهای فرزند به پردازش والد، نیازمند استفاده از Queue یا Pipe است.
۴. قطع کردن پردازشها
برخلاف threading، در multiprocessing میتوانید یک پردازش را با terminate() متوقف کنید. اما این کار خطرناک است و ممکن است منابع بهدرستی آزاد نشوند.
۵. خطاهای ارتباط بینپردازهای (IPC)
ارتباط بین پردازشها از طریق Queue، Pipe یا Shared Memory ممکن است با خطاهایی مانند BrokenPipeError یا ConnectionError مواجه شود که نیاز به مدیریت دارند.
الگوی Timeout در multiprocessing
multiprocessing روشهای مختلفی برای مدیریت Timeout ارائه میدهد. برخلاف asyncio که ابزارهای توکار قدرتمندی دارد، در multiprocessing باید از ترکیب روشهای مختلف استفاده کنیم.
روش اول: استفاده از Process.join(timeout)
سناریوی واقعی: یک سیستم پردازش تصویر دارید که باید یک تصویر بزرگ را با فیلترهای مختلف پردازش کند. هر فیلتر در یک پردازش جداگانه اجرا میشود. اگر پردازش یک فیلتر بیش از ۱۰ ثانیه طول بکشد، میخواهید آن را متوقف کنید و از نتیجهی فیلتر قبلی استفاده کنید.
import multiprocessing
import time
import random
def apply_filter(image_path, filter_name, result_queue):
"""شبیهسازی اعمال یک فیلتر روی تصویر"""
print(f"applying {filter_name} filter...")
time.sleep(random.uniform(2.0, 15.0))
if random.random() < 0.2:
raise ValueError(f"{filter_name} filter failed")
result_queue.put(f"{filter_name}_filtered_image")
def process_image_with_timeout(image_path, filters, timeout_per_filter=10.0):
"""پردازش تصویر با فیلترهای مختلف و Timeout"""
result_queue = multiprocessing.Queue()
processes = []
results = []
for filter_name in filters:
process = multiprocessing.Process(
target=apply_filter,
args=(image_path, filter_name, result_queue)
)
process.start()
processes.append((filter_name, process))
for filter_name, process in processes:
process.join(timeout_per_filter)
if process.is_alive():
print(f"timeout: {filter_name} filter took too long")
process.terminate()
process.join()
results.append(f"{filter_name}: timeout")
else:
try:
result = result_queue.get_nowait()
results.append(result)
except multiprocessing.queues.Empty:
print(f"error: {filter_name} filter produced no result")
results.append(f"{filter_name}: no result")
return results
if __name__ == "__main__":
filters = ["blur", "sharpen", "edge_detection", "sepia", "emboss"]
results = process_image_with_timeout("image.jpg", filters, timeout_per_filter=8.0)
print("\nfilter results:")
for result in results:
print(f" {result}")روش دوم: استفاده از concurrent.futures.ProcessPoolExecutor
سناریوی واقعی: یک سرویس تحلیل دادههای مالی دارید که باید چندین گزارش سنگین را همزمان تولید کند. هر گزارش در یک پردازش جداگانه تولید میشود. میخواهید اگر تولید یک گزارش بیش از ۵ ثانیه طول کشید، آن را لغو کنید و به کاربر پیام “تولید گزارش با تأخیر مواجه شد” نمایش دهید.
from concurrent.futures import ProcessPoolExecutor, TimeoutError
import time
import random
def generate_report(report_id, data_size):
"""شبیهسازی تولید یک گزارش مالی"""
print(f"generating report {report_id}...")
processing_time = random.uniform(0.5, 8.0)
time.sleep(processing_time)
if random.random() < 0.15:
raise ValueError(f"invalid data for report {report_id}")
return f"report_{report_id}_content_size_{data_size}mb"
def generate_reports_with_timeout(report_configs, timeout_per_report=5.0):
"""تولید چند گزارش با محدودیت زمان"""
results = {}
with ProcessPoolExecutor(max_workers=3) as executor:
futures = {
executor.submit(generate_report, config["id"], config["size"]): config["id"]
for config in report_configs
}
for future in futures:
report_id = futures[future]
try:
result = future.result(timeout=timeout_per_report)
results[report_id] = result
print(f"report {report_id} generated successfully")
except TimeoutError:
print(f"timeout generating report {report_id}")
results[report_id] = "timeout"
future.cancel()
except Exception as e:
print(f"error generating report {report_id}: {e}")
results[report_id] = f"error: {e}"
return results
if __name__ == "__main__":
reports = [
{"id": 1, "size": 10},
{"id": 2, "size": 20},
{"id": 3, "size": 15},
{"id": 4, "size": 30},
{"id": 5, "size": 25}
]
results = generate_reports_with_timeout(reports, timeout_per_report=4.0)
print("\nreport generation results:")
for report_id, result in results.items():
print(f" report {report_id}: {result}")روش سوم: استفاده از multiprocessing.Queue.get(timeout)
سناریوی واقعی: یک سیستم Crawler دارید که باید اطلاعات قیمت محصولات را از چندین سایت دریافت کند. هر سایت در یک پردازش جداگانه Crawl میشود. نتایج از طریق Queue به پردازش اصلی ارسال میشوند. اگر یک سایت بیش از ۳ ثانیه پاسخ ندهد، از آن صرف نظر میکنید.
import multiprocessing
import time
import random
def crawl_site(site_name, result_queue):
"""شبیهسازی Crawl کردن یک سایت"""
print(f"crawling {site_name}...")
delay = random.uniform(0.5, 5.0)
time.sleep(delay)
if random.random() < 0.25:
raise ConnectionError(f"cannot connect to {site_name}")
result_queue.put({site_name: f"data_from_{site_name}"})
def crawl_sites_with_timeout(sites, timeout_per_site=3.0):
"""Crawl کردن سایتها با Timeout با استفاده از Queue"""
result_queue = multiprocessing.Queue()
processes = []
results = {}
for site in sites:
process = multiprocessing.Process(
target=crawl_site,
args=(site, result_queue)
)
process.start()
processes.append(process)
# دریافت نتایج با Timeout
for _ in range(len(sites)):
try:
result = result_queue.get(timeout=timeout_per_site)
results.update(result)
print(f"received data from {list(result.keys())[0]}")
except multiprocessing.queues.Empty:
print(f"timeout waiting for site")
# توقف پردازشهای باقیمانده
for process in processes:
if process.is_alive():
print(f"terminating process {process.name}")
process.terminate()
process.join()
return results
if __name__ == "__main__":
sites = ["site1.com", "site2.com", "site3.com", "site4.com", "site5.com", "site6.com"]
results = crawl_sites_with_timeout(sites, timeout_per_site=3.0)
print(f"\nsuccessful crawls ({len(results)}):")
for site, data in results.items():
print(f" {site}: {data}")الگوی Retry در multiprocessing
پیادهسازی Retry در multiprocessing با توجه به سربار بالای ایجاد پردازش، نیازمند دقت بیشتری است. تعداد تلاشها باید کمتر و زمانهای انتظار طولانیتر باشند.
Retry با تأخیر ثابت (Fixed Delay)
سناریوی واقعی: یک سیستم تبدیل فایل دارید که فایلهای ویدیویی را به فرمتهای مختلف تبدیل میکند. گاهی فرآیند تبدیل بهدلیل کمبود منابع یا خطای موقت با شکست مواجه میشود. میخواهید تا ۲ بار با فاصلهی ۳ ثانیه دوباره تلاش کنید.
import multiprocessing
import time
import random
def convert_video(video_path, output_format, result_queue):
"""شبیهسازی تبدیل یک فایل ویدیویی"""
print(f"converting {video_path} to {output_format}...")
# شبیهسازی خطاهای موقت
if random.random() < 0.4:
raise MemoryError("insufficient memory for conversion")
time.sleep(random.uniform(0.5, 2.0))
result_queue.put(f"converted_{video_path}_to_{output_format}")
def retry_video_conversion(video_path, output_format, max_retries=2, delay=3.0):
"""تبدیل فایل با Retry"""
for attempt in range(max_retries):
result_queue = multiprocessing.Queue()
process = multiprocessing.Process(
target=convert_video,
args=(video_path, output_format, result_queue)
)
process.start()
process.join()
if process.exitcode != 0:
print(f"attempt {attempt + 1} for {video_path} failed with exit code {process.exitcode}")
if attempt == max_retries - 1:
raise Exception(f"all retry attempts failed for {video_path}")
print(f"waiting {delay} seconds before retry...")
time.sleep(delay)
else:
try:
result = result_queue.get_nowait()
print(f"conversion successful on attempt {attempt + 1}")
return result
except multiprocessing.queues.Empty:
print(f"attempt {attempt + 1} produced no result")
if attempt == max_retries - 1:
raise Exception(f"no result from conversion for {video_path}")
time.sleep(delay)
raise Exception("unreachable")
if __name__ == "__main__":
videos = ["movie1.mp4", "movie2.mp4", "movie3.mp4"]
for video in videos:
try:
result = retry_video_conversion(video, "avi", max_retries=2, delay=3.0)
print(f"result: {result}")
except Exception as e:
print(f"conversion failed for {video}: {e}")Retry با تأخیر تصاعدی (Exponential Backoff)
سناریوی واقعی: یک سرویس پردازش دادههای علمی دارید که محاسبات سنگین را روی مجموعهدادههای بزرگ انجام میدهد. گاهی سرور دادهها بهدلیل شلوغی پاسخ نمیدهد. با Exponential Backoff، در صورت خطا، مدت انتظار را افزایش میدهید تا به سرور فشار نیاورید.
import multiprocessing
import time
import random
def scientific_computation(dataset_id, result_queue):
"""شبیهسازی یک محاسبهی علمی سنگین"""
print(f"processing dataset {dataset_id}...")
if random.random() < 0.3:
raise ConnectionError(f"data server for dataset {dataset_id} unavailable")
time.sleep(random.uniform(1.0, 3.0))
result_queue.put(f"computed_result_for_{dataset_id}")
def retry_with_exponential_backoff(dataset_id, max_retries=3, base_delay=2.0):
"""پردازش داده با Exponential Backoff"""
for attempt in range(max_retries):
result_queue = multiprocessing.Queue()
process = multiprocessing.Process(
target=scientific_computation,
args=(dataset_id, result_queue)
)
process.start()
process.join()
if process.exitcode != 0:
print(f"attempt {attempt + 1} for dataset {dataset_id} failed")
if attempt == max_retries - 1:
raise Exception(f"all attempts failed for dataset {dataset_id}")
delay = base_delay * (2 ** attempt)
jitter = random.uniform(0, delay * 0.2)
total_delay = delay + jitter
print(f"waiting {total_delay:.2f} seconds before retry...")
time.sleep(total_delay)
else:
try:
result = result_queue.get_nowait()
print(f"dataset {dataset_id} processed successfully on attempt {attempt + 1}")
return result
except multiprocessing.queues.Empty:
print(f"attempt {attempt + 1} produced no result")
if attempt == max_retries - 1:
raise Exception(f"no result for dataset {dataset_id}")
time.sleep(base_delay * (2 ** attempt))
raise Exception("unreachable")
if __name__ == "__main__":
datasets = ["D001", "D002", "D003", "D004"]
for dataset in datasets:
try:
result = retry_with_exponential_backoff(dataset, max_retries=3, base_delay=2.0)
print(f"result: {result}")
except Exception as e:
print(f"processing failed for dataset {dataset}: {e}")ترکیب Timeout و Retry در multiprocessing
سناریوی واقعی: یک سرویس تحلیل ویدیو دارید که باید ویدیوهای آپلودشده را تحلیل کند. هر ویدیو در یک پردازش جداگانه تحلیل میشود. تحلیل ویدیو ممکن است زمانبر باشد یا با خطا مواجه شود. میخواهید تا ۲ بار تلاش کنید و هر بار حداکثر ۱۵ ثانیه صبر کنید.
import multiprocessing
import time
import random
def analyze_video(video_id, result_queue):
"""شبیهسازی تحلیل ویدیو"""
print(f"analyzing video {video_id}...")
duration = random.uniform(2.0, 20.0)
time.sleep(duration)
if random.random() < 0.25:
raise ValueError(f"video {video_id} is corrupted")
result_queue.put(f"analysis_result_for_{video_id}")
def analyze_with_timeout_and_retry(video_id, max_retries=2, timeout_per_attempt=15.0):
"""تحلیل ویدیو با Timeout و Retry ترکیبی"""
for attempt in range(max_retries):
print(f"\nattempt {attempt + 1} for video {video_id}")
result_queue = multiprocessing.Queue()
process = multiprocessing.Process(
target=analyze_video,
args=(video_id, result_queue)
)
process.start()
process.join(timeout=timeout_per_attempt)
if process.is_alive():
print(f"attempt {attempt + 1} timed out")
process.terminate()
process.join()
if attempt == max_retries - 1:
return {"video_id": video_id, "status": "timeout"}
delay = 2.0 * (2 ** attempt)
print(f"waiting {delay} seconds before retry...")
time.sleep(delay)
continue
if process.exitcode != 0:
print(f"attempt {attempt + 1} failed with exit code {process.exitcode}")
if attempt == max_retries - 1:
return {"video_id": video_id, "status": "error", "exitcode": process.exitcode}
delay = 2.0 * (2 ** attempt)
print(f"waiting {delay} seconds before retry...")
time.sleep(delay)
continue
try:
result = result_queue.get_nowait()
print(f"video {video_id} analyzed successfully")
return {"video_id": video_id, "status": "success", "result": result}
except multiprocessing.queues.Empty:
print(f"attempt {attempt + 1} produced no result")
if attempt == max_retries - 1:
return {"video_id": video_id, "status": "error", "error": "no result"}
return {"video_id": video_id, "status": "error", "error": "all attempts failed"}
if __name__ == "__main__":
videos = ["vid_001", "vid_002", "vid_003", "vid_004"]
results = []
for video in videos:
result = analyze_with_timeout_and_retry(video, max_retries=2, timeout_per_attempt=12.0)
results.append(result)
print("\n--- analysis results ---")
for result in results:
status = result["status"]
if status == "success":
print(f"video {result['video_id']}: success - {result['result']}")
elif status == "timeout":
print(f"video {result['video_id']}: timeout - analysis took too long")
else:
print(f"video {result['video_id']}: error - {result.get('error', 'unknown error')}")مدیریت استثناها در پردازشها
مدیریت استثناهایی که در پردازشهای فرزند رخ میدهند، یکی از چالشهای مهم در multiprocessing است. برخلاف threading، استثناهای پردازشها بهطور خودکار به پردازش والد منتقل نمیشوند.
روش اول: استفاده از Queue برای انتقال استثناها
سناریوی واقعی: یک سیستم پردازش فایل دارید که چندین فایل را همزمان پردازش میکند. هر فایل در یک پردازش جداگانه پردازش میشود. اگر پردازش یک فایل با خطا مواجه شود، میخواهید خطا را به پردازش والد گزارش دهید تا مشخص شود کدام فایل پردازش نشده است.
import multiprocessing
import time
import random
def process_file(file_path, result_queue):
"""شبیهسازی پردازش یک فایل"""
try:
print(f"processing {file_path}...")
time.sleep(random.uniform(0.5, 1.5))
if "corrupted" in file_path:
raise ValueError(f"file {file_path} is corrupted")
if "large" in file_path:
raise MemoryError(f"file {file_path} is too large")
if random.random() < 0.15:
raise PermissionError(f"permission denied for {file_path}")
result_queue.put(f"processed_{file_path}")
except Exception as e:
result_queue.put(e)
def batch_process_files(file_paths, max_workers=3):
"""پردازش دستهای فایلها با مدیریت خطا"""
result_queue = multiprocessing.Queue()
processes = []
results = {}
for path in file_paths:
process = multiprocessing.Process(
target=process_file,
args=(path, result_queue)
)
process.start()
processes.append(process)
for process in processes:
process.join()
while not result_queue.empty():
result = result_queue.get()
if isinstance(result, Exception):
error_msg = str(result)
print(f"error during processing: {error_msg}")
results[error_msg] = "failed"
else:
print(f"success: {result}")
results[result] = "success"
return results
if __name__ == "__main__":
files = [
"file1.txt",
"file2_corrupted.txt",
"file3.txt",
"file4_large.txt",
"file5.txt",
"file6.txt"
]
results = batch_process_files(files, max_workers=3)
print(f"\nprocessed {len([r for r in results.values() if r == 'success'])} files successfully")روش دوم: استفاده از Pipe برای انتقال استثناها
سناریوی واقعی: یک سرویس محاسبات ریاضی دارید که عملیاتهای سنگین را روی ماتریسهای بزرگ انجام میدهد. هر عملیات در یک پردازش جداگانه اجرا میشود. از Pipe برای ارتباط دوطرفه بین پردازش والد و فرزند استفاده میکنید.
import multiprocessing
import time
import random
def calculate_matrix(operation, matrix_size, conn):
"""شبیهسازی محاسبه روی ماتریس"""
try:
print(f"performing {operation} on {matrix_size}x{matrix_size} matrix...")
time.sleep(random.uniform(0.5, 2.0))
if operation == "inverse" and matrix_size > 100:
raise ValueError(f"matrix of size {matrix_size} is too large for inverse")
if random.random() < 0.2:
raise MemoryError("insufficient memory for operation")
conn.send(f"{operation}_result_{matrix_size}")
except Exception as e:
conn.send(e)
finally:
conn.close()
def run_operation(operation, matrix_size, timeout=5.0):
"""اجرای یک عملیات با Pipe و Timeout"""
parent_conn, child_conn = multiprocessing.Pipe()
process = multiprocessing.Process(
target=calculate_matrix,
args=(operation, matrix_size, child_conn)
)
process.start()
if parent_conn.poll(timeout):
result = parent_conn.recv()
process.join()
if isinstance(result, Exception):
raise result
return result
else:
process.terminate()
process.join()
raise TimeoutError(f"operation {operation} timed out")
if __name__ == "__main__":
operations = [
{"operation": "multiply", "size": 50},
{"operation": "inverse", "size": 150},
{"operation": "transpose", "size": 75},
{"operation": "determinant", "size": 100}
]
for op in operations:
try:
result = run_operation(op["operation"], op["size"], timeout=3.0)
print(f"{op['operation']} on {op['size']}x{op['size']}: {result}")
except TimeoutError as e:
print(f"timeout: {e}")
except ValueError as e:
print(f"value error: {e}")
except MemoryError as e:
print(f"memory error: {e}")
except Exception as e:
print(f"unexpected error: {e}")روش سوم: استفاده از concurrent.futures.ProcessPoolExecutor برای دریافت استثنا
سناریوی واقعی: یک سیستم تجزیهوتحلیل دادههای بزرگ دارید که چندین تحلیل آماری را همزمان روی مجموعهدادههای مختلف اجرا میکند. با استفاده از ProcessPoolExecutor، استثناها بهطور خودکار به پردازش والد منتقل میشوند.
from concurrent.futures import ProcessPoolExecutor
import time
import random
def statistical_analysis(dataset_id, analysis_type):
"""شبیهسازی تحلیل آماری روی یک مجموعهداده"""
print(f"analyzing dataset {dataset_id} with {analysis_type}...")
time.sleep(random.uniform(0.5, 2.0))
if dataset_id == "D003":
raise ValueError(f"dataset {dataset_id} contains invalid data")
if analysis_type == "regression" and random.random() < 0.3:
raise MemoryError("insufficient memory for regression analysis")
return f"{analysis_type}_result_for_{dataset_id}"
def run_analyses(analyses):
"""اجرای چندین تحلیل با مدیریت خطا"""
results = {}
with ProcessPoolExecutor(max_workers=3) as executor:
futures = {
executor.submit(statistical_analysis, a["dataset"], a["type"]): a
for a in analyses
}
for future in futures:
analysis = futures[future]
dataset = analysis["dataset"]
analysis_type = analysis["type"]
try:
result = future.result(timeout=4.0)
results[f"{dataset}_{analysis_type}"] = result
print(f"analysis {dataset}/{analysis_type} completed")
except TimeoutError:
print(f"timeout: analysis {dataset}/{analysis_type} took too long")
results[f"{dataset}_{analysis_type}"] = "timeout"
except ValueError as e:
print(f"value error for {dataset}/{analysis_type}: {e}")
results[f"{dataset}_{analysis_type}"] = f"error: {e}"
except MemoryError as e:
print(f"memory error for {dataset}/{analysis_type}: {e}")
results[f"{dataset}_{analysis_type}"] = f"error: {e}"
except Exception as e:
print(f"unexpected error for {dataset}/{analysis_type}: {e}")
results[f"{dataset}_{analysis_type}"] = f"error: {e}"
return results
if __name__ == "__main__":
analyses = [
{"dataset": "D001", "type": "mean"},
{"dataset": "D002", "type": "variance"},
{"dataset": "D003", "type": "regression"},
{"dataset": "D004", "type": "correlation"},
{"dataset": "D005", "type": "regression"},
{"dataset": "D006", "type": "mean"}
]
results = run_analyses(analyses)
print("\nanalysis results:")
for analysis_id, result in results.items():
print(f" {analysis_id}: {result}")مدیریت پردازشهای خراب و بازیابی منابع
سناریوی واقعی: یک سیستم پردازش دادههای جاری دارید که باید دادههای ورودی را در چندین مرحله پردازش کند. اگر یکی از پردازشها بهدلیل خطای غیرمنتظره از کار بیفتد، میخواهید آن را شناسایی کرده، منابع را پاکسازی کنید و دوباره راهاندازی کنید.
import multiprocessing
import time
import random
def data_processor(chunk_id, result_queue):
"""شبیهسازی پردازش یک تکه داده"""
print(f"processing chunk {chunk_id}...")
if random.random() < 0.3:
raise RuntimeError(f"unexpected error in chunk {chunk_id}")
time.sleep(random.uniform(0.3, 1.5))
result_queue.put(f"processed_chunk_{chunk_id}")
def process_with_recovery(chunks, max_retries=2, recovery_delay=2.0):
"""پردازش دادهها با قابلیت بازیابی"""
results = []
failed_chunks = []
result_queue = multiprocessing.Queue()
for chunk in chunks:
process = multiprocessing.Process(
target=data_processor,
args=(chunk, result_queue)
)
process.start()
process.join(timeout=3.0)
if process.is_alive():
print(f"chunk {chunk} timed out")
process.terminate()
process.join()
failed_chunks.append(chunk)
elif process.exitcode != 0:
print(f"chunk {chunk} failed with exit code {process.exitcode}")
failed_chunks.append(chunk)
else:
try:
result = result_queue.get_nowait()
results.append(result)
print(f"chunk {chunk} processed successfully")
except multiprocessing.queues.Empty:
print(f"chunk {chunk} produced no result")
failed_chunks.append(chunk)
# بازیابی تکههای شکستخورده
if failed_chunks and max_retries > 0:
print(f"\nrecovering {len(failed_chunks)} failed chunks...")
time.sleep(recovery_delay)
# بازیابی با پارامترهای جدید
recovered_results = process_with_recovery(
failed_chunks,
max_retries=max_retries - 1,
recovery_delay=recovery_delay * 1.5
)
results.extend(recovered_results)
return results
if __name__ == "__main__":
chunks = ["chunk_001", "chunk_002", "chunk_003", "chunk_004", "chunk_005"]
results = process_with_recovery(chunks, max_retries=2, recovery_delay=2.0)
print(f"\ncompleted {len(results)} out of {len(chunks)} chunks")
for result in results:
print(f" {result}")Idempotency در multiprocessing
سناریوی واقعی: یک سیستم پردازش تراکنش بانکی دارید که تراکنشها را در چندین پردازش موازی پردازش میکند. اگر یک تراکنش بهدلیل خطای شبکه با شکست مواجه شود، سیستم بهطور خودکار آن را دوباره اجرا میکند. باید مطمئن شوید که یک تراکنش بیش از یک بار پردازش نمیشود.
import multiprocessing
import time
import uuid
import random
class TransactionManager:
def __init__(self):
self.processed_transactions = set()
self.lock = multiprocessing.Lock()
def process_transaction(self, transaction_id, amount, idempotency_key, result_queue):
"""پردازش تراکنش با قابلیت Idempotency"""
try:
with self.lock:
if idempotency_key in self.processed_transactions:
print(f"transaction {transaction_id} already processed")
result_queue.put("already_processed")
return
print(f"processing transaction {transaction_id}...")
time.sleep(random.uniform(0.3, 1.0))
if amount > 1000 and random.random() < 0.3:
raise ValueError("amount exceeds limit")
self.processed_transactions.add(idempotency_key)
print(f"transaction {transaction_id} processed successfully")
result_queue.put("success")
except Exception as e:
result_queue.put(e)
def process_with_retry(manager, transaction_id, amount, max_retries=3):
"""پردازش تراکنش با Retry و Idempotency"""
idempotency_key = f"{transaction_id}_{str(uuid.uuid4())[:8]}"
for attempt in range(max_retries):
result_queue = multiprocessing.Queue()
process = multiprocessing.Process(
target=manager.process_transaction,
args=(transaction_id, amount, idempotency_key, result_queue)
)
process.start()
process.join()
if not result_queue.empty():
result = result_queue.get()
if isinstance(result, Exception):
print(f"attempt {attempt + 1} for {transaction_id} failed: {result}")
if attempt == max_retries - 1:
return "failed"
time.sleep(0.5 * (2 ** attempt))
else:
return result
else:
print(f"attempt {attempt + 1} for {transaction_id} produced no result")
if attempt == max_retries - 1:
return "failed"
time.sleep(0.5 * (2 ** attempt))
return "failed"
if __name__ == "__main__":
# برای multiprocessing، باید Manager را در main تعریف کرد
manager = TransactionManager()
transactions = [
{"id": "TXN-001", "amount": 500},
{"id": "TXN-002", "amount": 1200},
{"id": "TXN-003", "amount": 300},
{"id": "TXN-004", "amount": 800}
]
for tx in transactions:
result = process_with_retry(manager, tx["id"], tx["amount"], max_retries=3)
print(f"transaction {tx['id']}: {result}")مدیریت Pool و خطاهای آن
سناریوی واقعی: یک سیستم پردازش تصویر دارید که باید ۱۰۰ تصویر را با استفاده از multiprocessing.Pool پردازش کند. اگر پردازش یک تصویر با خطا مواجه شود، باید خطا را لاگ کنید و پردازش سایر تصاویر ادامه پیدا کند.
from multiprocessing import Pool
import time
import random
def process_image(image_id):
"""شبیهسازی پردازش یک تصویر"""
print(f"processing image {image_id}...")
time.sleep(random.uniform(0.1, 0.5))
if image_id % 7 == 0:
raise ValueError(f"image {image_id} is corrupted")
if image_id % 13 == 0:
raise MemoryError(f"image {image_id} is too large")
return f"processed_image_{image_id}"
def process_images_with_pool(image_ids, num_workers=4):
"""پردازش تصاویر با Pool و مدیریت خطا"""
results = {}
errors = []
with Pool(processes=num_workers) as pool:
async_results = []
for image_id in image_ids:
async_result = pool.apply_async(process_image, (image_id,))
async_results.append((image_id, async_result))
for image_id, async_result in async_results:
try:
result = async_result.get(timeout=3.0)
results[image_id] = result
print(f"image {image_id} processed successfully")
except TimeoutError:
print(f"timeout processing image {image_id}")
errors.append(f"image {image_id}: timeout")
except ValueError as e:
print(f"value error for image {image_id}: {e}")
errors.append(f"image {image_id}: {e}")
except MemoryError as e:
print(f"memory error for image {image_id}: {e}")
errors.append(f"image {image_id}: {e}")
except Exception as e:
print(f"unexpected error for image {image_id}: {e}")
errors.append(f"image {image_id}: {e}")
return results, errors
if __name__ == "__main__":
image_ids = list(range(1, 101))
results, errors = process_images_with_pool(image_ids, num_workers=4)
print(f"\nprocessed {len(results)} images successfully")
print(f"failed {len(errors)} images")
if errors:
print("\nsample errors:")
for error in errors[:5]:
print(f" {error}")بهترین شیوهها در مدیریت خطای multiprocessing
۱. همیشه از join(timeout) برای مدیریت زمان انتظار استفاده کنید
این روش به شما اجازه میدهد پردازشهایی که بیش از حد طول میکشند را شناسایی و متوقف کنید.
۲. از terminate() با احتیاط استفاده کنید
terminate() پردازش را بهصورت ناگهانی متوقف میکند و ممکن است منابع بهدرستی آزاد نشوند. تا حد امکان از روشهای هماهنگی برای توقف پردازشها استفاده کنید.
۳. برای انتقال استثناها از Queue یا Pipe استفاده کنید
استثناهای پردازشهای فرزند بهطور خودکار به پردازش والد منتقل نمیشوند. برای دریافت خطاها از Queue یا Pipe استفاده کنید.
۴. تعداد Retry را محدود کنید
به دلیل سربار بالای ایجاد پردازشها، تعداد تلاشهای مجدد را کمتر از threading یا asyncio در نظر بگیرید (معمولاً ۲ تا ۳ بار کافی است).
۵. از Exponential Backoff با Jitter استفاده کنید
این کار از هجوم به سرویسهای خارجی و سیستمهای دیگر جلوگیری میکند.
۶. از apply_async با مدیریت خطا در Pool استفاده کنید
Pool.apply_async به شما اجازه میدهد خطاهای هر کار را بهصورت جداگانه مدیریت کنید.
۷. از ProcessPoolExecutor برای مدیریت آسانتر استفاده کنید
این API مدیریت Timeout و استثناها را سادهتر میکند.
۸. اطمینان از Idempotency قبل از استفاده از Retry
اگر عملیات Idempotent نیست، برای Retry کردن باید از شناسههای یکتا (Idempotency Key) استفاده کنید.
۹. لاگهای دقیق نگه دارید
در هر مرحله از Timeout و Retry، لاگ بنویسید تا بتوانید مشکلات را ردیابی کنید.
۱۰. از if __name__ == "__main__" استفاده کنید
در سیستمهای ویندوز، استفاده از این شرط برای جلوگیری از ایجاد حلقههای بینهایت در ایجاد پردازشها ضروری است.
جمعبندی
مدیریت خطا در multiprocessing نیازمند رویکردی متفاوت از threading و asyncio است. سربار بالای ایجاد پردازشها، حافظهی جداگانه و نیاز به ارتباط بینپردازهای، همه عواملی هستند که باید در طراحی مدیریت خطا در نظر گرفته شوند.
الگوهای اصلی که در این فصل بررسی کردیم عبارتند از:
Timeout: با استفاده از Process.join(timeout)، ProcessPoolExecutor با result(timeout) و Queue.get(timeout)
Retry: با تأخیر ثابت، تأخیر تصاعدی و ترکیب با Timeout
مدیریت استثناها: با استفاده از Queue، Pipe و ProcessPoolExecutor
Idempotency: برای اطمینان از ایمنی Retry کردن
مدیریت Pool: با apply_async و مدیریت خطاهای هر کار
با رعایت این اصول و الگوها، میتوانید برنامههای multiprocessing بنویسید که در برابر خطاهای موقت، طولانی شدن عملیاتها، خرابی پردازشها و شرایط غیرمنتظره، مقاوم باشند. به خاطر داشته باشید که multiprocessing برای عملیاتهای CPU-bound طراحی شده است و اگر پروژهی شما I/O-bound است، asyncio یا threading انتخابهای بهتری هستند.