مدیریت خطا در multiprocessing

Please login to bookmark Close

در این فصل، به‌صورت کامل و مستقل به بررسی مدیریت خطا در محیط 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 انتخاب‌های بهتری هستند.

Please login to bookmark Close
پیشرفت شما در «دوره آموزش کانکارنسی در پایتون» (91%)
نظرات

دیدگاهتان را بنویسید

91%
پیشرفت

سرفصل دوره

فهرست مطالب

سرفصل دوره

تمرین

این قسمت تمرین ندارد!

پاسخ تمرین ها

هنوز برای تمرین‌های این قسمت پاسخی ثبت نشده است!

اشتراک گذاری

چرا بهتره از فیلترشکن استفاده کنید؟

من همه ویدئو ها و پادکست های کُدباز رو توی یوتیوب و ساندکلود و پلتفرم هایی آپلود می‌کنم که اغلب فیلتر هستند.

اغلب آموزش‌ها ویدئو و پادکست دارند. پس اگر می‌خواهید از محتوای سایت بیشترین استفاده رو ببرید نیاز به فیلتر شکن دارید.

توجه داشته باشید که برای خرید از فروشگاه بهتره فیلتر شکن رو خاموش کنید.

تنظیمات

انتخاب زبان
تغییر تم