Il pool multiprocessing Python si arresta bruscamente
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
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()
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()