Docker本地运行Airflow DAG部分任务持续处于scheduled状态求助
问题解决:Airflow本地Docker环境下任务卡在"Scheduled"状态,Pool提示负可用槽位
问题核心
本地Docker运行Airflow时,指定使用dbrelay资源池的任务始终处于scheduled状态,实例详情提示:
('Not scheduling since there are %s open slots in pool %s and require %s pool slots', -1, 'dbrelay', 1)
小型DAG偶尔能正常运行,包含大量并发任务的大型DAG必现该问题,生产环境无此异常。
根因分析
生产环境中已提前创建并配置好dbrelay资源池,但本地Docker环境默认不会自动创建自定义Pool。当任务指定使用不存在的Pool时,Airflow会默认给该Pool分配0个总槽位,此时请求该Pool的任务会计算出负的可用槽位(已占用0,需求1,可用=0-0-1=-1),导致任务无法被调度。
解决方案
1. 手动创建dbrelay资源池(推荐)
进入本地Airflow UI的Admin > Pools页面:
- 点击"Create"新建Pool
- Pool Name填写
dbrelay - Slots根据本地资源设置(建议和生产环境保持一致,比如5或10)
- Description可填写"DB Relay任务专用池"
2. 临时测试方案
如果仅为快速验证代码,可临时修改任务代码,去掉pool="dbrelay"参数,让任务使用默认的default_pool:
op_dbrelay = EPBigQueryInsertJobOperator( task_id=f"db_relay_{relay}", sql="templates/generic/create_or_replace_snapshot.sql", # 注释或删除pool参数 # pool="dbrelay", params={ "SOURCE_PROJECT": PROJECT, "DESTINATION_PROJECT": PROJECT, "SOURCE_DATASET": "dbrelay", "DESTINATION_DATASET": TMP_VIEW_DATASET, "TABLE": relay, }, )
3. Docker环境自动初始化Pool
如果需要每次启动Docker都自动创建Pool,可以在Airflow的$AIRFLOW_HOME/dags目录下添加初始化脚本,使用Airflow Python API自动创建:
from airflow.models import Pool from airflow.utils.session import create_session def create_dbrelay_pool(): with create_session() as session: pool = session.query(Pool).filter(Pool.pool == "dbrelay").first() if not pool: new_pool = Pool( pool="dbrelay", slots=5, # 按需设置槽位数 description="DB Relay tasks pool" ) session.add(new_pool) session.commit() create_dbrelay_pool()
启动Airflow时该脚本会自动执行,完成Pool的创建。
附用户提供的相关代码
default_args = { "owner": "airflow", "depends_on_past": False, "start_date": datetime(2022, 2, 18), "email": ["xxx@xxx.de"], "retries": 1, "retry_delay": timedelta(minutes=5), "on_failure_callback": task_fail_slack_alert, # slack alert on all task failures "location": "eu", "bigquery_conn_id": BIGQUERY_CONN_ID, "provide_context": True, } with DAG( "dag_name", default_args=default_args, concurrency=3, # concurrent tasks catchup=False, schedule_interval="30 05,17 * * *", max_active_runs=1, params={ "PROJECT": PROJECT, "TMP_VIEW_DATASET": TMP_VIEW_DATASET, "HOURS_TO_SUBTRACT": 4, "REPORT_DATE": REPORT_DATE, }, ) as dag: with TaskGroup(group_id="dbrelays") as tg_dbrelays: for index, relay in enumerate(dbrelays): op_dbrelay = EPBigQueryInsertJobOperator( task_id=f"db_relay_{relay}", sql="templates/generic/create_or_replace_snapshot.sql", pool="dbrelay", params={ "SOURCE_PROJECT": PROJECT, "DESTINATION_PROJECT": PROJECT, "SOURCE_DATASET": "dbrelay", "DESTINATION_DATASET": TMP_VIEW_DATASET, "TABLE": relay, }, )
内容的提问来源于stack exchange,提问作者Alex
相关产品推荐
相关产品推荐

