Estoy tratando de acceder a archivos externos en una tarea de flujo de aire para leer algunos sql, y obtengo "archivo no encontrado". ¿Se ha topado alguien con esto?
from airflow import DAG from airflow.operators.python_operator import PythonOperator from datetime import datetime, timedelta dag = DAG( 'my_dat', start_date=datetime(2017, 1, 1), catchup=False, schedule_interval=timedelta(days=1) ) def run_query(): # read the query query = open('sql/queryfile.sql') # run the query execute(query) tas = PythonOperator( task_id='run_query', dag=dag, python_callable=run_query)El registro indica lo siguiente:
IOError: [Errno 2] No such file or directory: 'sql/queryfile.sql'Entiendo que podría simplemente copiar y pegar la consulta dentro del mismo archivo, realmente no es una solución ordenada. Hay múltiples consultas y el texto es realmente grande, incrustarlo con el código de Python comprometería la legibilidad.
Aquí hay un ejemplo de uso de Variable para hacerlo más fácil.
Primero agregue Variable en Airflow UI -> Admin -> Variable , por ejemplo. {key: 'sql_path', values: 'your_sql_script_folder'}
Luego agregue el siguiente código en su DAG, para usar la Variable de Airflow que acaba de agregar.
Código DAG:
import airflow from airflow.models import Variable tmpl_search_path = Variable.get("sql_path") dag = airflow.DAG( 'tutorial', schedule_interval="@daily", template_searchpath=tmpl_search_path, # this default_args=default_args )Ahora puede usar el nombre o la ruta del script sql en la carpeta Variable
Puedes aprender más en este
Todas las rutas relativas se toman en referencia a la variable de entorno AIRFLOW_HOME . Tratar:
Suponiendo que el directorio sql es relativo al archivo Python actual, puede averiguar la ruta absoluta al archivo sql de esta manera:
import os CUR_DIR = os.path.abspath(os.path.dirname(__file__)) def run_query(): # read the query query = open(f"{CUR_DIR}/sql/queryfile.sql") # run the query execute(query)