Estoy creando una clase de administrador de subprocesos que maneja la ejecución de tareas como subprocesos y pasa los resultados al siguiente paso del proceso. El flujo funciona correctamente en la primera ejecución de recibir una tarea, pero la segunda ejecución falla con el siguiente error:
...python3.8/concurrent/futures/thread.py", line 179, in submit raise RuntimeError('cannot schedule new futures after shutdown') RuntimeError: cannot schedule new futures after shutdown Las tareas provienen de la entrada del usuario de Cmd.cmdloop ; por lo tanto, el script es persistente y está destinado a no cerrarse. En su lugar, run se llamará varias veces, a medida que se reciba la entrada del usuario.
Implementé un ThreadPoolExecutor para manejar la carga de trabajo y tratar de recopilar los resultados cronológicamente con concurrent.futures.as_completed para que cada elemento se procese al siguiente paso en orden de finalización.
El método de run a continuación funciona perfectamente para la primera ejecución, pero devuelve el error en la segunda ejecución de la misma tarea (que tuvo éxito durante la primera ejecución).
def run ( self, _executor=None, _futures={}, ) -> bool : task = self.pipeline.get( ) with _executor or self.__default_executor as executor : _futures = { executor.submit ( task.target.execute, ), } for future in concurrent.futures.as_completed ( _futures, ) : print( future.result ( ) ) return True Entonces, la idea es que cada llamada a run creará y desmantelará el executor con el contexto. Pero el error sugiere que el contexto se cerró correctamente después de la primera ejecución, y no se puede volver a abrir/recrear cuando se llama a run durante la segunda iteración... ¿a qué apunta este error? .. ¿Qué me estoy perdiendo?
Cualquier ayuda sería genial, gracias de antemano.
Su solución más fácil será usar la biblioteca de multiprocesamiento en lugar de enviar futuros y ThreadPoolExecutore con Context Manager:
pool = ThreadPool(50) pool.starmap(test_function, zip(array1,array2...)) pool.close() pool.join() Mientras que (array1[0] , array2[0]) serán los valores enviados a la función "test_function" en el primer subproceso, (array1[1] , array2[1]) en el segundo subproceso, y así sucesivamente.