Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

426
Views
Reinicie periódicamente el grupo de multiprocesamiento de Python

Tengo un grupo de multiprocesamiento de Python que realiza un trabajo muy largo que, incluso después de una depuración exhaustiva, no es lo suficientemente sólido como para no fallar cada 24 horas aproximadamente, porque depende de muchas herramientas de terceros que no son de Python con interacciones complejas. Además, la máquina subyacente tiene ciertos problemas que no puedo controlar. Tenga en cuenta que al fallar no me refiero a que todo el programa se bloquee, sino que algunos o la mayoría de los procesos se vuelven inactivos debido a algunos errores, y la aplicación en sí se bloquea o continúa el trabajo solo con los procesos que no han fallado.

Mi solución en este momento es matar periódicamente el trabajo, manualmente, y luego simplemente reiniciar desde donde estaba.

Aunque no sea lo ideal, lo que quiero hacer ahora es lo siguiente: reiniciar el grupo de multiprocesamiento periódicamente, programáticamente, desde el propio código de Python. Realmente no me importa si esto implica matar a los trabajadores de la piscina en medio de su trabajo. ¿Cuál sería la mejor manera de hacerlo?

Mi código se parece a:

 with Pool() as p: for _ in p.imap_unordered(function, data): save_checkpoint() log()

Lo que tengo en mente sería algo como:

 start = 0 end = 1000 # magic number while start + 1 < len(data): current_data = data[start:end] with Pool() as p: for _ in p.imap_unordered(function, current_data): save_checkpoint() log() start += 1 end += 1

O:

 start = 0 end = 1000 # magic number while start + 1 < len(data): current_data = data[start:end] start_timeout(time=TIMEOUT) # which would be the best way to to do that without breaking multiprocessing? try: with Pool() as p: for _ in p.imap_unordered(function, current_data): save_checkpoint() log() start += 1 end += 1 except Timeout: pass

O cualquier sugerencia que creas que sería mejor. Cualquier ayuda sería muy apreciada, gracias!

over 4 years ago · Santiago Trujillo
1 answers
Answer question

0

El problema con su código actual es que itera los resultados multiprocesados directamente y esa llamada se bloqueará. Afortunadamente, hay una solución fácil: use apply_async exactamente como se sugiere en los documentos . Pero debido a cómo describe el caso de uso aquí y la falla, lo he adaptado un poco. En primer lugar, una tarea simulada:

 from multiprocessing import Pool, TimeoutError, cpu_count from time import sleep from random import randint def log(): print("logging is a dangerous activity: wear a hard hat.") def work(d): sleep(randint(1, 100) / 100) print("finished working") if randint(1, 10) == 1: print("blocking...") while True: sleep(0.1) return d

Esta función de trabajo fallará con una probabilidad de 0.1 , bloqueándose indefinidamente. Creamos las tareas:

 data = list(range(100)) nproc = cpu_count()

Y luego generar futuros para todos ellos:

 while data: print(f"== Processing {len(data)} items. ==") with Pool(nproc) as p: tasks = [p.apply_async(work, (d,)) for d in data]

Entonces podemos intentar sacar las tareas manualmente:

 for task in tasks: try: res = task.get(timeout=1) data.remove(res) log() except TimeoutError: failed.append(task) if len(failed) < nproc: print( f"{len(failed)} processes are blocked," f" but {nproc - len(failed)} remain." ) else: break

El tiempo de espera de control aquí es el tiempo de espera para .get . Debe ser tan largo como espera que tome el proceso más largo. Tenga en cuenta que detectamos cuando todo el grupo está atado y nos damos por vencidos.

Pero dado que en el escenario que describe, algunos subprocesos tardarán más que otros, podemos dar a los procesos 'fallidos' algo de tiempo para recuperarse. Así, cada vez que falla una tarea, comprobamos rápidamente si las demás han tenido éxito:

 for task in failed: try: res = task.get(timeout=0.01) data.remove(res) failed.remove(task) log() except TimeoutError: continue

Si esta es una buena adición en su caso depende de si sus tareas son realmente tan inestables como supongo que son.

Salir del administrador de contexto para el grupo terminará el grupo, por lo que ni siquiera necesitamos manejar eso nosotros mismos. Si tiene una variación significativa, es posible que desee aumentar el tamaño del grupo (aumentando así la cantidad de tareas que pueden detenerse) o permitir que las tareas tengan un período de gracia antes de considerarlas "fallidas".

over 4 years ago · Santiago Trujillo Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!