Otras preguntas sobre "tareas dinámicas" parecen abordar la construcción dinámica de un DAG en el momento del cronograma o del diseño. Estoy interesado en agregar dinámicamente tareas a un DAG durante la ejecución.
from airflow import DAG from airflow.operators.dummy_operator import DummyOperator from airflow.operators.python_operator import PythonOperator from datetime import datetime dag = DAG('test_dag', description='a test', schedule_interval='0 0 * * *', start_date=datetime(2018, 1, 1), catchup=False) def make_tasks(): du1 = DummyOperator(task_id='dummy1', dag=dag) du2 = DummyOperator(task_id='dummy2', dag=dag) du3 = DummyOperator(task_id='dummy3', dag=dag) du1 >> du2 >> du3 p = PythonOperator( task_id='python_operator', dag=dag, python_callable=make_tasks)Esta implementación ingenua no parece funcionar: las tareas ficticias nunca aparecen en la interfaz de usuario.
¿Cuál es la forma correcta de agregar nuevos operadores al DAG durante la ejecución? ¿Es posible?
No es posible modificar el DAG durante su ejecución (sin mucho más trabajo).
El programador recoge dag = DAG(... en un bucle. Tendrá la instancia de tarea 'python_operator' . Esa instancia de tarea se programa en una ejecución de dag y la ejecuta un trabajador o ejecutor. Dado que los modelos DAG en Airflow DB solo las actualiza el programador, estas tareas ficticias agregadas no se conservarán en el DAG ni se programarán para ejecutarse. Se olvidarán cuando el trabajador salga. A menos que copie todo el código del programador con respecto a la persistencia y actualización del modelo. … pero eso se deshará la próxima vez que el programador visite el archivo DAG para analizarlo, lo que podría ocurrir una vez por minuto, una vez por segundo o más rápido dependiendo de cuántos otros archivos DAG haya para analizar.
Airflow en realidad quiere que cada DAG mantenga aproximadamente el mismo diseño entre ejecuciones. También quiere recargar/analizar archivos DAG constantemente. Entonces, aunque podría crear un archivo DAG que en cada ejecución determine las tareas dinámicamente en función de algunos datos externos (preferiblemente almacenados en caché en un archivo o módulo pyc, no E/S de red como una búsqueda de base de datos, reducirá la velocidad de todo el ciclo de programación para todos los DAG) no es un buen plan, ya que su gráfico y su vista de árbol se volverán confusos, y el análisis de su programador se verá más afectado por esa búsqueda.
Podrías hacer que el invocable ejecute cada tarea...
def make_tasks(context): du1 = DummyOperator(task_id='dummy1', dag=dag) du2 = DummyOperator(task_id='dummy2', dag=dag) du3 = DummyOperator(task_id='dummy3', dag=dag) du1.execute(context) du2.execute(context) du3.execute(context) p = PythonOperator( provides_context=true,Pero eso es secuencial, y debe descubrir cómo usar Python para hacerlos paralelos (¿usar futuros?) Y si alguno genera una excepción, toda la tarea falla. También está vinculado a un ejecutor o trabajador, por lo que no utiliza la distribución de tareas de airflow (kubernetes, mesos, celery).
La otra forma de trabajar con esto es agregar un número fijo de tareas (el número máximo) y usar las llamadas para cortocircuitar las tareas innecesarias o enviar argumentos con xcom para cada una de ellas, cambiando su comportamiento en tiempo de ejecución. pero sin cambiar el DAG.
Con respecto a su muestra de código, nunca llama a su función que registra sus tareas en su DAG.
Para tener una especie de tareas dinámicas, puede tener un solo operador que haga algo diferente dependiendo de algún estado o puede tener un puñado de operadores que se pueden omitir dependiendo del estado, con un ShortCircuitOperator.
Aprecio todo el trabajo que todos han hecho aquí, ya que tengo el mismo desafío de crear DAG dinámicamente estructurados. He cometido suficientes errores como para no usar el software en contra de su diseño. Si no puedo inspeccionar toda la ejecución en la interfaz de usuario y acercar y alejar, básicamente uso las funciones de flujo de aire, que son la razón principal por la que lo uso de todos modos. Puedo simplemente escribir código de multiprocesamiento dentro de una función y terminar con eso también.
Habiendo dicho todo eso, mi solución es usar un administrador de recursos como el bloqueo de redis y tener un DAG que escribe en este administrador de recursos con datos sobre qué ejecutar, cómo ejecutar, etc. y tener otro DAG o DAG que se ejecuten en ciertos intervalos sondeando al administrador de recursos, bloqueándolos antes de ejecutarlos y eliminándolos al finalizar. De esta manera, al menos uso el flujo de aire como se esperaba, aunque sus especificaciones no satisfacen exactamente mis necesidades. Desgloso el problema en partes más definibles. Las soluciones son creativas pero van en contra del diseño y no han sido probadas por los desarrolladores. Dicen específicamente tener flujos de trabajo estructurados fijos. No puedo poner un trabajo alrededor del código que no está probado y contra el diseño a menos que reescriba el código de flujo de aire central y lo pruebe yo mismo. Entiendo que mi solución trae complejidad con el bloqueo y todo eso, pero al menos conozco los límites de eso.