He usado Joblib y Airflow en el pasado y no me he encontrado con este problema. Estoy tratando de ejecutar un trabajo a través de Airflow que ejecuta un cálculo paralelo usando Joblib. Cuando se inicia el trabajo de Airflow, veo la siguiente advertencia
UserWarning: Loky-backed parallel loops cannot be called in multiprocessing, setting n_jobs=1Al rastrear la advertencia hasta el origen, veo que se activa la siguiente función en el paquete joblib en la clase LokyBackend (también hay una lógica similar en la clase MultiprocessingBackend)
def effective_n_jobs(self, n_jobs): """Determine the number of jobs which are going to run in parallel""" if n_jobs == 0: raise ValueError('n_jobs == 0 in Parallel has no meaning') elif mp is None or n_jobs is None: # multiprocessing is not available or disabled, fallback # to sequential mode return 1 elif mp.current_process().daemon: # Daemonic processes cannot have children if n_jobs != 1: warnings.warn( 'Loky-backed parallel loops cannot be called in a' ' multiprocessing, setting n_jobs=1', stacklevel=3) return 1 El problema es que he ejecutado una función similar en Joblib y Airflow antes y no activé esta condición para establecer n_jobs en 1. Me pregunto si se trata de algún tipo de problema de control de versiones (usando Airflow 2.X y Joblib 1.X). ) o si hay configuraciones en Airflow que pueden arreglar esto. Miré versiones anteriores de Joblib e incluso cambié a Joblib 0.4.0, pero eso no resolvió ningún problema. Dudo más en degradar Airflow debido a las diferencias en la API, las conexiones de la base de datos, etc.
Editar:
Aquí está el código que he estado ejecutando en Airflow:
def test_parallel(): out=joblib.Parallel(n_jobs=-1, backend="loky")(joblib.delayed(lambda a: a+1)(i) for i in range(20)) with DAG("test", default_args=DEFAULT_ARGS, schedule_interval="0 8 * * *",) as test: run_test = PythonOperator( task_id="test", python_callable=test_parallel, ) run_testY la salida en los registros de flujo de aire:
[2021-07-27 10:41:29,890] {logging_mixin.py:104} WARNING - /data01/code/virtualenv/alpha/lib/python3.8/site-packages/joblib/parallel.py:733 UserWarning: Loky-backed parallel loops cannot be called in a multiprocessing, setting n_jobs=1 Lanzo el airflow scheduler aire y el airflow webserver de flujo de aire a través de supervisor . Sin embargo, incluso si ejecuto ambos procesos de flujo de aire desde la línea de comandos, el problema persiste. Sin embargo, no sucede cuando solo ejecuto la tarea a través de la API de tareas de flujo de aire, por ejemplo airflow tasks test run_test
Noté que no llamaste a la función run_test en la parte inferior de tu código. ¿Podría ser esa la causa de algún problema? Versión corregida:
def test_parallel(): out=joblib.Parallel(n_jobs=-1, backend="loky")(joblib.delayed(lambda a: a+1)(i) for i in range(20)) with DAG("test", default_args=DEFAULT_ARGS, schedule_interval="0 8 * * *",) as test: run_test = PythonOperator( task_id="test", python_callable=test_parallel, ) run_test()Así que resolví esto cambiando de PythonOperator a BashOpertaor y joblib se detuvo para reducir el número de CPU y subprocesos a 1. También seguí las instrucciones de aquí solo para eliminar los procesos demoníacos después de la ejecución del código, pero puede esperar 300 segundos, que es el tiempo de espera predeterminado de joblib. para los procesos que terminan.