Il pool multiprocessing Python si arresta bruscamente

Sep 27 2020

Sto cercando di eseguire l'elaborazione parallela per le mie esigenze e il codice sembra funzionare come previsto per elementi 4k-5k in parallelo. Ma non appena gli elementi da elaborare iniziano ad aumentare, il codice elabora alcuni elenchi e quindi, senza generare alcun errore, il programma si interrompe bruscamente.

Ho controllato e il programma non si blocca, la RAM è disponibile (ho una RAM da 16 Gb) e l'utilizzo della CPU non è nemmeno del 30%. Non riesco a capire cosa sta succedendo. Ho 1 milione di elementi da elaborare.

def get_items_to_download():
    #iterator to fetch all items that are to be downloaded
    yield download_item

def start_download_process():
    multiproc_pool = multiprocessing.Pool(processes=10)
    for download_item in get_items_to_download():
        multiproc_pool.apply_async(start_processing, args = (download_item, ), callback = results_callback)
    
    multiproc_pool.close()
    multiproc_pool.join()

def start_processing(download_item):
    try:
        # Code to download item from web API
        # Code to perform some processing on the data
        # Code to update data into database
        return True
    except Exception as e:
        return False

def results_callback(result):
    print(result)

if __name__ == "__main__":
    start_download_process()

AGGIORNARE -

Trovato l'errore - BrokenPipeError: [Errno 32] tubo rotto

Traccia -

Traceback (most recent call last):
File "/usr/lib/python3.6/multiprocessing/pool.py", line 125, in worker
put((job, i, result))
File "/usr/lib/python3.6/multiprocessing/queues.py", line 347, in put
self._writer.send_bytes(obj)
File "/usr/lib/python3.6/multiprocessing/connection.py", line 200, in send_bytes
self._send_bytes(m[offset:offset + size])
File "/usr/lib/python3.6/multiprocessing/connection.py", line 404, in _send_bytes
self._send(header + buf)
File "/usr/lib/python3.6/multiprocessing/connection.py", line 368, in _send
n = write(self._handle, buf)
BrokenPipeError: [Errno 32] Broken pipe

Risposte

Simplecode Oct 10 2020 at 16:04
def get_items_to_download():
    #instead of yield, return the complete generator object to avoid iterating over this function.
    #Return type - generator (download_item1, download_item2...)
    return download_item


def start_download_process():
    download_item = get_items_to_download()
    # specify the chunksize to get faster results. 
    with multiprocessing.Pool(processes=10) as pool:
    #map_async() is also available, if that's your use case.
        results= pool.map(start_processing, download_item, chunksize=XX )  
    print(results)
    return(results)

def start_processing(download_item):
    try:
        # Code to download item from web API
        # Code to perform some processing on the data
        # Code to update data into database
        return True
    except Exception as e:
        return False

def results_callback(result):
    print(result)

if __name__ == "__main__":
    start_download_process()
Booboo Oct 04 2020 at 16:30

Il codice sembra corretto. L'unica cosa a cui riesco a pensare è che tutti i tuoi processi sono sospesi in attesa del completamento. Ecco un suggerimento: invece di utilizzare il meccanismo di callback fornito da apply_async, utilizzare l' AsyncResultoggetto restituito per ottenere il valore restituito dal processo. Puoi chiamaregetsu questo oggetto specificando un valore di timeout (30 secondi arbitrariamente specificati di seguito, possibilmente non abbastanza lunghi). Se l'attività non è stata completata in quella durata, verrà generata un'eccezione di timeout (puoi prenderla, se lo desideri). Ma questo metterà alla prova l'ipotesi che i processi siano sospesi. Assicurati solo di specificare un valore di timeout sufficientemente grande da completare l'attività entro quel periodo di tempo. Ho anche suddiviso gli invii di attività in batch di 1000, non perché penso che la dimensione di 1.000.000 sia un problema di per sé , ma solo così non hai un elenco di 1.000.000 di oggetti risultato. Ma se ti accorgi che non ti blocchi più, prova ad aumentare la dimensione del batch e vedi se fa la differenza.

import multiprocessing

def get_items_to_download():
    #iterator to fetch all items that are to be downloaded
    yield download_item

BATCH_SIZE = 1000

def start_download_process():
    with multiprocessing.Pool(processes=10) as multiproc_pool:
        results = []
        for download_item in get_items_to_download():
            results.append(multiproc_pool.apply_async(start_processing, args = (download_item, )))
            if len(results) == BATCH_SIZE:
                process_results(results)
                results = []
        if len(results):
            process_results(results)
    

def start_processing(download_item):
    try:
        # Code to download item from web API
        # Code to perform some processing on the data
        # Code to update data into database
        return True
    except Exception as e:
        return False

TIMEOUT_VALUE = 30 # or some suitable value

def process_results(results):
    for result in results:
        return_value = result.get(TIMEOUT_VALUE) # will cause an exception if process is hanging
        print(return_value)

if __name__ == "__main__":
    start_download_process()

Aggiornare

Sulla base della ricerca su Google di diverse pagine per l'errore di pipe rotto, sembra che l'errore potrebbe essere il risultato di esaurire la memoria. Vedere Multiprocessing Python: eccezione del tubo rotto dopo aver aumentato la dimensione del pool , per esempio. La seguente rielaborazione tenta di utilizzare meno memoria. Se funziona, puoi provare ad aumentare la dimensione del batch:

import multiprocessing


BATCH_SIZE = 1000
POOL_SIZE = 10


def get_items_to_download():
    #iterator to fetch all items that are to be downloaded
    yield download_item


def start_download_process():
    with multiprocessing.Pool(processes=POOL_SIZE) as multiproc_pool:
        items = []
        for download_item in get_items_to_download():
            items.append(download_item)
            if len(items) == BATCH_SIZE:
                process_items(multiproc_pool, items)
                items = []
        if len(items):
            process_items(multiproc_pool, items)


def start_processing(download_item):
    try:
        # Code to download item from web API
        # Code to perform some processing on the data
        # Code to update data into database
        return True
    except Exception as e:
        return False


def compute_chunksize(iterable_size):
    if iterable_size == 0:
        return 0
    chunksize, extra = divmod(iterable_size, POOL_SIZE * 4)
    if extra:
        chunksize += 1
    return chunksize


def process_items(multiproc_pool, items):
    chunksize = compute_chunksize(len(items))
    # you must iterate the iterable returned:
    for return_value in multiproc_pool.imap(start_processing, items, chunksize):
        print(return_value)


if __name__ == "__main__":
    start_download_process()