Docker化Airflow DAG执行两次问题排查求助
Airflow DAG Docker部署后自动运行两次问题排查与解决
我将实现CSV文件转存至GCP、加载至BigQuery并完成简单转换的Airflow DAG进行Docker化部署后,遇到一个问题:执行docker-compose run时DAG会自动运行两次,但通过Airflow UI手动触发时仅运行一次,无法定位问题根源。
我的DAG代码
from scripts import extract_and_gcpload, load_to_BQ from datetime import datetime default_args = { 'owner': 'shweta', 'start_date': datetime(2025, 4, 24), 'retries': 0 } with DAG( 'spacex_etl_dag', default_args=default_args, schedule_interval=None, schedule=None, catchup=False # 阻止Airflow运行错过的周期任务 ) as dag: extract_and_upload = PythonOperator( task_id="extract_and_upload_to_gcs", python_callable=extract_and_gcpload.load_to_gcp_pipeline, ) load_to_bq = PythonOperator( task_id="load_to_BQ", python_callable=load_to_BQ.load_csv_to_bigquery ) run_dbt = BashOperator( task_id="run_dbt", bash_command="cd '/opt/airflow/dbt/my_dbt' && dbt run --profiles-dir /opt/airflow/dbt" ) extract_and_upload >> load_to_bq >> run_dbt
入口文件startscript.sh
#!/bin/bash set -euo pipefail log() { echo "[$(date +'%Y-%m-%d %H:%M:%S')] $1" } # 可选:运行数据库初始化并解析DAG log "Initializing Airflow DB..." airflow db upgrade log "Parsing DAGs..." airflow scheduler --num-runs 1 DAG_ID="spacex_etl_dag" log "Unpausing DAG: $DAG_ID" airflow dags unpause "$DAG_ID" || true log "Triggering DAG: $DAG_ID" airflow dags trigger "$DAG_ID" || true log "Creating admin user (if not exists)..." airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email admin@example.com \ --password admin || true if [[ "$1" == "webserver" || "$1" == "scheduler" ]]; then log "Starting Airflow: $1" exec airflow "$@" else log "Executing: $@" exec "$@" fi
docker-compose.yaml文件
services: airflow-webserver: build: context: . dockerfile: Dockerfile container_name: airflow-webserver env_file: .env restart: always environment: AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'false' AIRFLOW__LOGGING__REMOTE_LOGGING: 'False' AIRFLOW__CORE__LOAD_EXAMPLES: 'false' AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow GOOGLE_APPLICATION_CREDENTIALS: /opt/airflow/secrets/llms-395417-c18ea70a3f54.json volumes: - ./dags:/opt/airflow/dags - ./scripts:/opt/airflow/scripts - ./dbt:/opt/airflow/dbt - ./secrets:/opt/airflow/secrets ports: - 8080:8080 command: webserver airflow-scheduler: build: context: . dockerfile: Dockerfile container_name: airflow-scheduler env_file: .env restart: always environment: AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'false' AIRFLOW__LOGGING__REMOTE_LOGGING: 'False' AIRFLOW__CORE__LOAD_EXAMPLES: 'false' AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow GOOGLE_APPLICATION_CREDENTIALS: /opt/airflow/secrets/llms-395417-c18ea70a3f54.json volumes: - ./dags:/opt/airflow/dags - ./dbt:/opt/airflow/dbt - ./secrets:/opt/airflow/secrets - ./scripts:/opt/airflow/scripts depends_on: - postgres command: scheduler postgres: image: postgres:13 environment: POSTGRES_USER: airflow POSTGRES_PASSWORD: airflow POSTGRES_DB: airflow volumes: - postgres-db-volume:/var/lib/postgresql/data volumes: postgres-db-volume:
问题根源分析
- 重复触发DAG:
airflow-webserver和airflow-scheduler两个容器基于同一镜像构建,镜像的入口脚本为startscript.sh。两个容器启动时都会执行脚本中的airflow dags trigger "spacex_etl_dag"命令,导致DAG被触发两次。 - 额外的调度触发:脚本中的
airflow scheduler --num-runs 1会启动一次scheduler并执行一轮任务调度,即使DAG设置了schedule_interval=None,该命令也可能意外触发DAG运行,叠加手动trigger的次数。
解决方案
方案1:拆分初始化步骤到独立容器
修改docker-compose.yaml,新增专门的初始化容器执行一次性操作,避免webserver和scheduler重复执行:
services: airflow-init: build: context: . dockerfile: Dockerfile container_name: airflow-init env_file: .env environment: AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow volumes: - ./dags:/opt/airflow/dags - ./scripts:/opt/airflow/scripts - ./dbt:/opt/airflow/dbt - ./secrets:/opt/airflow/secrets depends_on: - postgres command: > bash -c " airflow db upgrade && airflow users create --username admin --firstname Admin --lastname User --role Admin --email admin@example.com --password admin || true && airflow dags unpause spacex_etl_dag || true # 如需自动触发DAG,仅在这里执行一次 # airflow dags trigger spacex_etl_dag || true " # 保留原webserver和scheduler配置,无需修改它们的命令
方案2:修改startscript.sh,避免重复执行
修改脚本,仅在特定容器中执行初始化和trigger操作:
#!/bin/bash set -euo pipefail log() { echo "[$(date +'%Y-%m-%d %H:%M:%S')] $1" } # 仅在scheduler容器中执行初始化操作 if [[ "$1" == "scheduler" ]]; then log "Initializing Airflow DB..." airflow db upgrade log "Creating admin user (if not exists)..." airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email admin@example.com \ --password admin || true DAG_ID="spacex_etl_dag" log "Unpausing DAG: $DAG_ID" airflow dags unpause "$DAG_ID" || true # 如需自动触发DAG,仅在这里执行一次 # log "Triggering DAG: $DAG_ID" # airflow dags trigger "$DAG_ID" || true fi # 移除airflow scheduler --num-runs 1,避免额外调度 if [[ "$1" == "webserver" || "$1" == "scheduler" ]]; then log "Starting Airflow: $1" exec airflow "$@" else log "Executing: $@" exec "$@" fi
关键修改点
- 移除
startscript.sh中的airflow scheduler --num-runs 1命令,避免不必要的调度触发。 - 确保
airflow dags trigger命令仅执行一次,要么放在独立初始化容器中,要么只在一个服务容器中执行。 - 初始化操作(数据库升级、用户创建、DAG解锁)仅执行一次,禁止多个容器重复执行。
内容的提问来源于stack exchange,提问作者Shweta Dalal
相关产品推荐
相关产品推荐

