Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

446
Vistas
Python: concurrent.futures ¿Cómo hacerlo cancelable?

Python concurrent.futures y ProcessPoolExecutor proporcionan una interfaz ordenada para programar y monitorear tareas. Los futuros incluso proporcionan un método .cancel():

cancel() : intento de cancelar la llamada. Si la llamada se está ejecutando actualmente y no se puede cancelar , el método devolverá False; de lo contrario, la llamada se cancelará y el método devolverá True.

Desafortunadamente, en una pregunta similar (sobre asyncio), la respuesta afirma que las tareas en ejecución no se pueden cancelar con este recorte de la documentación, pero los documentos no dicen eso, solo si se están ejecutando Y no se pueden cancelar.

Enviar multiprocessing.Events a los procesos tampoco es trivialmente posible (hacerlo a través de parámetros como en multiprocess.Process devuelve un RuntimeError)

¿Qué estoy tratando de hacer? Me gustaría dividir un espacio de búsqueda y ejecutar una tarea para cada partición. Pero es suficiente tener UNA solución y el proceso requiere un uso intensivo de la CPU. Entonces, ¿hay una forma cómoda de lograr esto que no compense las ganancias al usar ProcessPool para empezar?

Ejemplo:

 from concurrent.futures import ProcessPoolExecutor, FIRST_COMPLETED, wait # function that profits from partitioned search space def m_run(partition): for elem in partition: if elem == 135135515: return elem return False futures = [] # used to create the partitions steps = 100000000 with ProcessPoolExecutor(max_workers=4) as pool: for i in range(4): # run 4 tasks with a partition, but only *one* solution is needed partition = range(i*steps,(i+1)*steps) futures.append(pool.submit(m_run, partition)) done, not_done = wait(futures, return_when=FIRST_COMPLETED) for d in done: print(d.result()) print("---") for d in not_done: # will return false for Cancel and Result for all futures print("Cancel: "+str(d.cancel())) print("Result: "+str(d.result()))
over 4 years ago · Santiago Trujillo
3 Respuestas
Responde la pregunta

0

Desafortunadamente, la ejecución de Futures no se puede cancelar. Creo que la razón principal es garantizar la misma API en diferentes implementaciones (no es posible interrumpir subprocesos o rutinas en ejecución).

La biblioteca Pebble fue diseñada para superar esta y otras limitaciones.

 from pebble import ProcessPool def function(foo, bar=0): return foo + bar with ProcessPool() as pool: future = pool.schedule(function, args=[1]) # if running, the container process will be terminated # a new process will be started consuming the next task future.cancel()
over 4 years ago · Santiago Trujillo Denunciar

0

No sé por qué concurrent.futures.Future no tiene un método .kill() , pero puede lograr lo que quiere cerrando el grupo de procesos con pool.shutdown(wait=False) y eliminando los procesos secundarios restantes manualmente.

Cree una función para matar procesos secundarios:

 import signal, psutil def kill_child_processes(parent_pid, sig=signal.SIGTERM): try: parent = psutil.Process(parent_pid) except psutil.NoSuchProcess: return children = parent.children(recursive=True) for process in children: process.send_signal(sig)

Ejecute su código hasta que obtenga el primer resultado, luego elimine todos los procesos secundarios restantes:

 from concurrent.futures import ProcessPoolExecutor, FIRST_COMPLETED, wait # function that profits from partitioned search space def m_run(partition): for elem in partition: if elem == 135135515: return elem return False futures = [] # used to create the partitions steps = 100000000 pool = ProcessPoolExecutor(max_workers=4) for i in range(4): # run 4 tasks with a partition, but only *one* solution is needed partition = range(i*steps,(i+1)*steps) futures.append(pool.submit(m_run, partition)) done, not_done = wait(futures, timeout=3600, return_when=FIRST_COMPLETED) # Shut down pool pool.shutdown(wait=False) # Kill remaining child processes kill_child_processes(os.getpid())
over 4 years ago · Santiago Trujillo Denunciar

0

Encontré su pregunta interesante, así que aquí está mi hallazgo.

Encontré que el comportamiento del método .cancel() es como se indica en la documentación de python. En cuanto a sus funciones concurrentes en ejecución, lamentablemente no pudieron cancelarse incluso después de que se les indicó que lo hicieran. Si mi hallazgo es correcto, entonces razono que Python requiere un método .cancel() más efectivo.

Ejecute el siguiente código para verificar mi hallazgo.

 from concurrent.futures import ProcessPoolExecutor, as_completed from time import time # function that profits from partitioned search space def m_run(partition): for elem in partition: if elem == 3351355150: return elem break #Added to terminate loop once found return False start = time() futures = [] # used to create the partitions steps = 1000000000 with ProcessPoolExecutor(max_workers=4) as pool: for i in range(4): # run 4 tasks with a partition, but only *one* solution is needed partition = range(i*steps,(i+1)*steps) futures.append(pool.submit(m_run, partition)) ### New Code: Start ### for f in as_completed(futures): print(f.result()) if f.result(): print('break') break for f in futures: print(f, 'running?',f.running()) if f.running(): f.cancel() print('Cancelled? ',f.cancelled()) print('New Instruction Ended at = ', time()-start ) print('Total Compute Time = ', time()-start )

Actualización: es posible finalizar por la fuerza los procesos concurrentes a través de bash, pero la consecuencia es que el programa principal de python también finalizará. Si esto no es un problema para usted, intente con el siguiente código.

Debe agregar los códigos a continuación entre las últimas 2 declaraciones de impresión para ver esto por sí mismo. Nota: este código solo funciona si no está ejecutando ningún otro programa python3.

 import subprocess, os, signal result = subprocess.run(['ps', '-C', 'python3', '-o', 'pid='], stdout=subprocess.PIPE).stdout.decode('utf-8').split() print ('result =', result) for i in result: print('PID = ', i) if i != result[0]: os.kill(int(i), signal.SIGKILL) try: os.kill(int(i), 0) raise Exception("""wasn't able to kill the process HINT:use signal.SIGKILL or signal.SIGABORT""") except OSError as ex: continue
over 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda