Airflow 1.10.6迁移至1.10.15时ExternalTaskSensor执行日期函数报TypeError错误求助
TypeError with ExternalTaskSensor execution_date_fn after upgrading to Airflow 1.10.15
这个问题是因为Airflow 1.10.10及后续版本对ExternalTaskSensor的execution_date_fn参数做了签名变更,直接导致你的旧代码参数不匹配。
问题根源
在Airflow 1.10.6中,execution_date_fn的函数签名只需要接收execution_date和自定义参数;但从1.10.10开始,这个函数会优先接收一个完整的上下文字典(context)作为第一个参数,而不是直接传入execution_date对象。你的get_updated_exec_date函数现在拿到的第一个参数是context字典,而非预期的datetime,后续执行execution_date + relativedelta以及传递delta时自然会抛出类型错误。
修复方案
你需要修改函数来从context中提取execution_date,同时如果要传递自定义参数(比如你的delta),可以用functools.partial来包装:
1. 更新日期处理函数
from airflow.utils.dates import relativedelta from datetime import timedelta def get_updated_exec_date(context, delta=-1): # 从context字典中取出execution_date execution_date = context['execution_date'] date_next_month = execution_date + relativedelta(months=1) start_next_month = date_next_month.replace(day=1) return start_next_month + timedelta(days=delta)
2. 修改Sensor的调用方式
用partial来传递delta参数,确保context能被正确传入:
from functools import partial sensor = ExternalTaskSensor( task_id='wait_for_parent', external_dag_id='parent', external_task_id='reporting', allowed_states=['success'], execution_date_fn=partial(get_updated_exec_date, delta=-1), dag=dag )
如果你的delta参数固定不需要动态调整,也可以直接简化函数,去掉delta参数:
def get_updated_exec_date(context): execution_date = context['execution_date'] date_next_month = execution_date + relativedelta(months=1) start_next_month = date_next_month.replace(day=1) return start_next_month + timedelta(days=-1)
这样修改后,代码就能在Airflow 1.10.15中正常运行了。
内容的提问来源于stack exchange,提问作者Cam
相关产品推荐
相关产品推荐

