Actualmente estoy tratando de usar Airflow para orquestar un proceso en el que algunos operadores se definen dinámicamente y dependen de la salida de otro operador (anterior).
En el código a continuación, t1 actualiza un archivo de texto con nuevos registros (estos en realidad se leen desde una cola externa pero, para simplificar, los codifiqué como A, B y C aquí). Luego, quiero crear operadores separados para cada registro leído de ese archivo de texto. Estos operadores crearán los directorios A, B y C, respectivamente, y en la interfaz de usuario de Airflow se verán como procesos bash independientes Create_directory_A, Create_directory_B y Create_directory_C.
dag = DAG('Test_DAG', description="Lorem ipsum.", start_date=datetime(2017, 3, 20), schedule_interval=None, catchup=False) def create_text_file(list_of_rows): text_file = open('text_file.txt', "w") for row in list_of_rows: text_file.write(row + '\n') text_file.close() def read_text(): txt_file = open('text_file.txt', 'r') return [element for element in txt_file.readlines()] t1 = PythonOperator( task_id='Create_text_file', python_callable=create_text_file, op_args=[['A', 'B', 'C']], dag=dag ) for row in read_text(): t2 = BashOperator( task_id='Create_directory_{}'.format(row), bash_command="mkdir {{params.dir_name}}", params={'dir_name': row}, dag=dag ) t1 >> t2En la documentación de Airflow, puedo ver que el programador lo ejecutará [DAG] periódicamente para reflejar los cambios, si los hay . ¿Eso significa que existe el riesgo de que, aunque mi operador t1 se ejecute antes que t2, los operadores bash se creen para la lista de registros antes de la actualización (ya que fue entonces cuando se evaluó el DAG)?
No puede crear tareas dinámicamente que dependan del resultado de una tarea anterior. Estás confundiendo el cronograma y el tiempo de ejecución. Se crea una definición de DAG y una tarea a la hora programada. Se crea una instancia de ejecución y tarea de DAG en el momento de la ejecución. Solo una instancia de tarea puede generar resultados.
El programador de Airflow creará el gráfico dinámico con lo que contenga text_file.txt en el momento del programa . Estas tareas luego se envían a los trabajadores.
Un trabajador eventualmente ejecutará la instancia de la tarea t1 y creará un nuevo text_file.txt , pero en este punto, el programador ya calculó la lista de tareas t2 y la envió a los trabajadores.
Por lo tanto, cualquiera que sea la instancia de tarea t1 más reciente que se descargue en text_file.txt , se usará la próxima vez que el programador decida que es hora de ejecutar el DAG.
Si su tarea es rápida y sus trabajadores no están atrasados, ese será el contenido de la ejecución DAG anterior. Si están atrasados, el contenido de text_file.txt puede estar obsoleto y, si no tiene mucha suerte, el programador lee el archivo mientras una instancia de tarea lo está escribiendo y obtendrá datos incompletos de read_text() .
Este código en realidad creará una instancia de t2 que será construida por el operador bash con la última row que obtiene de read_text() . Estoy seguro de que esto no es lo que quieres.
Un mejor enfoque sería crear un DAG separado para su operador t2 que se activa cuando t1 escribe el archivo. Hay una pregunta SO sobre esto que podría ayudar: Apache Airflow: desencadenar/programar la repetición de DAG al finalizar (Sensor de archivos)