Solucionar problemas de DAGs

Cloud Composer 3 | Cloud Composer 2 | Cloud Composer 1

En esta página se ofrecen pasos para solucionar problemas habituales de los flujos de trabajo e información sobre ellos.

Algunos problemas de ejecución de DAG pueden deberse a que el programador de Airflow no funciona correctamente o de forma óptima. Sigue las instrucciones para solucionar problemas del programador para resolverlos.

Solucionar problemas de flujo de trabajo

Para empezar a solucionar el problema, haz lo siguiente:

  1. Consulta los registros de Airflow.

    Puedes aumentar el nivel de registro de Airflow anulando la siguiente opción de configuración de Airflow.

    Sección Clave Valor
    logging (core en Airflow 1) logging_level El valor predeterminado es INFO. Define el valor DEBUG para obtener más detalles en los mensajes de registro.
  2. Consulta el panel de control de monitorización.

  3. Consulta Cloud Monitoring.

  4. En la Google Cloud consola, comprueba si hay errores en las páginas de los componentes de tu entorno.

  5. En la interfaz web de Airflow, consulta la vista de gráfico del DAG para ver las instancias de tareas fallidas.

    Sección Clave Valor
    webserver dag_orientation LR, TB, RL o BT

Depurar errores del operador

Para depurar un error del operador, sigue estos pasos:

  1. Comprueba si hay errores específicos de la tarea.
  2. Consulta los registros de Airflow.
  3. Consulta Cloud Monitoring.
  4. Consulta los registros específicos del operador.
  5. Corrige los errores.
  6. Sube el DAG a la carpeta /dags.
  7. En la interfaz web de Airflow, borra los estados anteriores del DAG.
  8. Reanuda o ejecuta el DAG.

Solucionar problemas de ejecución de tareas

Airflow es un sistema distribuido con muchas entidades, como el programador, el ejecutor y los trabajadores, que se comunican entre sí a través de una cola de tareas y la base de datos de Airflow, y envían señales (como SIGTERM). En el siguiente diagrama se muestra un resumen de las interconexiones entre los componentes de Airflow.

Interacción entre los componentes de Airflow
Figura 1. Interacción entre los componentes de Airflow (haz clic para ampliar)

En un sistema distribuido como Airflow, puede haber problemas de conectividad de red o la infraestructura subyacente puede experimentar problemas intermitentes. Esto puede provocar que las tareas fallen y se reprogramen para su ejecución, o que no se completen correctamente (por ejemplo, tareas zombi o tareas que se han quedado bloqueadas durante la ejecución). Airflow tiene mecanismos para hacer frente a estas situaciones y reanudar automáticamente el funcionamiento normal. En las siguientes secciones se explican los problemas habituales que se producen durante la ejecución de tareas de Airflow.

Las tareas fallan sin emitir ningún registro

La tarea falla sin emitir registros debido a errores de análisis de DAG

A veces, puede haber errores sutiles en los DAGs que provoquen que el programador de Airflow pueda programar tareas para su ejecución, que el procesador de DAGs pueda analizar el archivo de DAG, pero que el trabajador de Airflow no pueda ejecutar las tareas del DAG porque hay errores de programación en el archivo de DAG. Esto puede provocar que una tarea de Airflow se marque como Failed y no haya ningún registro de su ejecución.

Soluciones:

  • Verifica en los registros de los trabajadores de Airflow que no haya errores relacionados con la falta de un DAG o con errores de análisis de DAG.

  • Aumenta los parámetros relacionados con el análisis de DAG:

    • Aumenta [dagbag-import-timeout][ext-airflow-dagrun-import-timeout] a al menos 120 segundos (o más, si es necesario).

    • Aumenta el valor de dag-file-processor-timeout a 180 segundos como mínimo (o más, si es necesario). Este valor debe ser superior a dagbag-import-timeout.

  • Consulta también Solucionar problemas del procesador de DAG.

Las tareas se interrumpen de forma brusca

Durante la ejecución de las tareas, los workers de Airflow pueden finalizar de forma abrupta debido a problemas que no están relacionados específicamente con la tarea en sí. Consulta las causas raíz habituales para ver una lista de estos casos y las posibles soluciones. En las siguientes secciones se describen algunos síntomas adicionales que podrían derivarse de esas causas principales:

Tareas zombi

Airflow detecta dos tipos de discrepancias entre una tarea y un proceso que ejecuta la tarea:

  • Las tareas fallidas son tareas que deberían estar en ejecución, pero no lo están. Esto puede ocurrir si el proceso de la tarea se ha terminado o no responde, si el trabajador de Airflow no ha informado del estado de la tarea a tiempo porque está sobrecargado o si se ha apagado la máquina virtual en la que se ejecuta la tarea. Airflow busca estas tareas periódicamente y las rechaza o las vuelve a intentar, según la configuración de la tarea.

    Descubrir tareas zombi

    resource.type="cloud_composer_environment"
    resource.labels.environment_name="ENVIRONMENT_NAME"
    log_id("airflow-scheduler")
    textPayload:"Detected zombie job"
  • Las tareas zombi son tareas que no deberían estar en ejecución. Airflow busca estas tareas periódicamente y las finaliza.

Consulta las causas habituales para obtener más información sobre cómo solucionar problemas de tareas zombi.

Señales SIGTERM

Linux, Kubernetes, el programador de Airflow y Celery usan señales SIGTERM para finalizar los procesos responsables de ejecutar los trabajadores o las tareas de Airflow.

Puede haber varios motivos por los que se envían señales SIGTERM en un entorno:

  • Una tarea se ha convertido en una tarea fallida y debe detenerse.

  • El programador ha detectado un duplicado de una tarea y envía señales de instancia de finalización y SIGTERM a la tarea para detenerla.

  • En el autoescalado horizontal de pods, el plano de control de GKE envía señales SIGTERM para eliminar los pods que ya no son necesarios.

  • El programador puede enviar señales SIGTERM al proceso DagFileProcessorManager. El programador usa estas señales SIGTERM para gestionar el ciclo de vida del proceso DagFileProcessorManager y se pueden ignorar sin problemas.

    Ejemplo:

    Launched DagFileProcessorManager with pid: 353002
    Sending Signals.SIGTERM to group 353002. PIDs of all processes in the group: []
    Sending the signal Signals.SIGTERM to group 353002
    Sending the signal Signals.SIGTERM to process 353002 as process group is missing.
    
  • Condición de carrera entre la retrollamada de latido y las retrollamadas de salida en local_task_job, que monitoriza la ejecución de la tarea. Si el latido detecta que una tarea se ha marcado como correcta, no puede distinguir si la tarea en sí se ha completado correctamente o si se le ha indicado a Airflow que la considere correcta. Sin embargo, finalizará un ejecutor de tareas sin esperar a que se cierre.

    Estas señales SIGTERM se pueden ignorar sin problemas. La tarea ya está en estado correcto y la ejecución del DAG no se verá afectada.

    La entrada de registro Received SIGTERM. es la única diferencia entre la salida normal y la finalización de la tarea en el estado correcto.

    Condición de carrera entre las devoluciones de llamada de latido y de salida
    Imagen 2. Condición de carrera entre las devoluciones de llamada de la señal de latido y de salida (haz clic para ampliar)
  • Un componente de Airflow usa más recursos (CPU, memoria) de los permitidos por el nodo del clúster.

  • El servicio de GKE realiza operaciones de mantenimiento y envía señales SIGTERM a los pods que se ejecutan en un nodo que está a punto de actualizarse.

    Cuando se termina una instancia de tarea con SIGTERM, puedes ver las siguientes entradas de registro en los registros de un trabajador de Airflow que ha ejecutado la tarea:

    {local_task_job.py:211} WARNING - State of this instance has been externally
    set to queued. Terminating instance. {taskinstance.py:1411} ERROR - Received
    SIGTERM. Terminating subprocesses. {taskinstance.py:1703} ERROR - Task failed
    with exception
    

Posibles soluciones:

Este problema se produce cuando una máquina virtual que ejecuta la tarea se queda sin memoria. Esto no está relacionado con las configuraciones de Airflow, sino con la cantidad de memoria disponible para la VM.

  • En Cloud Composer 1, puedes volver a crear tu entorno con un tipo de máquina que tenga un mayor rendimiento.

  • Puedes reducir el valor de la opción de configuración de [celery]worker_concurrencyconcurrencia de Airflow. Esta opción determina cuántas tareas ejecuta simultáneamente un trabajador de Airflow.

Negsignal.SIGKILL ha interrumpido la tarea de Airflow

A veces, es posible que tu tarea use más memoria de la que se ha asignado al trabajador de Airflow. En ese caso, podría interrumpirse a las Negsignal.SIGKILL. El sistema envía esta señal para evitar un mayor consumo de memoria que podría afectar a la ejecución de otras tareas de Airflow. En el registro del trabajador de Airflow, puede que veas la siguiente entrada:

{local_task_job.py:102} INFO - Task exited with return code Negsignal.SIGKILL

Negsignal.SIGKILL también puede aparecer como código -9.

Posibles soluciones:

  • Menor worker_concurrency de trabajadores de Airflow.

  • Cambia a un tipo de máquina más grande que se use en el clúster de Cloud Composer.

  • Optimiza tus tareas para que usen menos memoria.

La tarea falla debido a la presión de los recursos

Síntoma: durante la ejecución de una tarea, el subproceso del trabajador de Airflow responsable de la ejecución de la tarea de Airflow se interrumpe de forma abrupta. El error que se muestra en el registro del trabajador de Airflow puede ser similar al siguiente:

...
File "/opt/python3.8/lib/python3.8/site-packages/celery/app/trace.py", line 412, in trace_task    R = retval = fun(*args, **kwargs)  File "/opt/python3.8/lib/python3.8/site-packages/celery/app/trace.py", line 704, in __protected_call__    return self.run(*args, **kwargs)  File "/opt/python3.8/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 88, in execute_command    _execute_in_fork(command_to_exec)  File "/opt/python3.8/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 99, in _execute_in_fork
raise AirflowException('Celery command failed on host: ' + get_hostname())airflow.exceptions.AirflowException: Celery command failed on host: airflow-worker-9qg9x
...

Solución:

La tarea falla debido a la expulsión de un pod

Los pods de Google Kubernetes Engine están sujetos al ciclo de vida de los pods de Kubernetes y al desalojo de pods. Los picos de tareas y la programación conjunta de trabajadores son dos de las causas más habituales de la expulsión de pods en Cloud Composer.

La expulsión de pods puede producirse cuando un pod concreto usa en exceso los recursos de un nodo en relación con las expectativas de consumo de recursos configuradas para el nodo. Por ejemplo, el desalojo puede producirse cuando se ejecutan varias tareas que consumen mucha memoria en un pod y su carga combinada hace que el nodo en el que se ejecuta este pod supere el límite de consumo de memoria.

Si se expulsa un pod de trabajador de Airflow, se interrumpirán todas las instancias de tareas que se estén ejecutando en ese pod y, más adelante, Airflow las marcará como fallidas.

Los registros se almacenan en búfer. Si se expulsa un pod de trabajador antes de que se vacíe el búfer, no se emitirán registros. Si una tarea falla y no hay registros, significa que los trabajadores de Airflow se han reiniciado debido a que se ha quedado sin memoria (OOM). Es posible que algunos registros estén presentes en Cloud Logging aunque no se hayan emitido los registros de Airflow.

Para ver los registros, sigue estos pasos:

  1. En la Google Cloud consola, ve a la página Entornos.

    Ir a Entornos

  2. En la lista de entornos, haz clic en el nombre del entorno. Se abrirá la página Detalles del entorno.

  3. Ve a la pestaña Registros.

  4. Para ver los registros de los distintos workers de Airflow, vaya a Todos los registros > Registros de Airflow > Workers.

Síntoma:

  1. En la Google Cloud consola, ve a la página Cargas de trabajo.

    Ir a Cargas de trabajo

  2. Si hay airflow-worker pods que muestran Evicted, haz clic en cada pod expulsado y busca el mensaje The node was low on resource: memory en la parte superior de la ventana.

Solución:

  • Crea un entorno de Cloud Composer 1 con un tipo de máquina más grande que el actual.

  • Consulta los registros de los pods de airflow-worker para ver los posibles motivos de desalojo. Para obtener más información sobre cómo obtener registros de Pods concretos, consulta el artículo Solucionar problemas con cargas de trabajo implementadas.

  • Asegúrate de que las tareas del DAG sean idempotentes y se puedan volver a intentar.

  • Evita descargar archivos innecesarios en el sistema de archivos local de los trabajadores de Airflow.

    Los workers de Airflow tienen una capacidad limitada en el sistema de archivos local. Cuando se agota el espacio de almacenamiento, el plano de control de GKE expulsa el pod de trabajo de Airflow. Se produce un error en todas las tareas que estaba ejecutando el trabajador expulsado.

    Ejemplos de operaciones problemáticas:

    • Descargar archivos u objetos y almacenarlos localmente en un trabajador de Airflow. En su lugar, almacena estos objetos directamente en un servicio adecuado, como un segmento de Cloud Storage.
    • Acceder a objetos grandes en la carpeta /data desde un trabajador de Airflow. El trabajador de Airflow descarga el objeto en su sistema de archivos local. En su lugar, implementa tus DAGs de forma que los archivos de gran tamaño se procesen fuera del pod de trabajador de Airflow.

Causas habituales

El trabajador de Airflow se ha quedado sin memoria

Cada trabajador de Airflow puede ejecutar hasta [celery]worker_concurrency instancias de tareas simultáneamente. Si el consumo acumulativo de memoria de esas instancias de tareas supera el límite de memoria de un trabajador de Airflow, se termina un proceso aleatorio en él para liberar recursos.

Descubrir eventos de falta de memoria de los trabajadores de Airflow

resource.type="k8s_node"
resource.labels.cluster_name="GKE_CLUSTER_NAME"
log_id("events")
jsonPayload.message:"Killed process"
jsonPayload.message:("airflow task" OR "celeryd")

A veces, la falta de memoria en un trabajador de Airflow puede provocar que se envíen paquetes mal formados durante una sesión de SQL Alchemy a la base de datos, a un servidor DNS o a cualquier otro servicio al que llame un DAG. En este caso, el otro extremo de la conexión podría rechazar o interrumpir las conexiones del trabajador de Airflow. Por ejemplo:

"UNKNOWN:Error received from peer
{created_time:"2024-