Airflow Sensor无法访问上下文变量问题排查及解决方案咨询
问题
我正在尝试构建一个Airflow Sensor,用于读取触发DAG时可配置的参数(触发DAG时可通过config修改),以此确定等待时长。以下是我的实现代码:
from airflow.decorators import dag, task, task_group from datetime import date, datetime, timedelta import re params = { "time":"8h" } def parse_time(time_str): regex = re.compile(r'^((?P<days>[\.\d]+?)d)?((?P<hours>[\.\d]+?)h)?((?P<minutes>[\.\d]+?)m)?((?P<seconds>[\.\d]+?)s)?$') parts = regex.match(time_str) if parts is None: return timedelta() time_params = {name: float(param) for name, param in parts.groupdict().items() if param} return timedelta(**time_params) @dag( dag_id="test", start_date=datetime(2023, 5, 1), schedule_interval = None, catchup=False, default_args={"retries":0}, params=params, tags=["test","debug"], ) def test(): @task.sensor( task_id=f"run_after", poke_interval=60 * 5, timeout=60 * 60 * 24 * 3, mode="reschedule" ) def run_after(**context): run_after = context["params"].get("time","0h") print(run_after) target_time = parse_time(run_after) time_since_midnight = datetime.now() - datetime.strptime(context["data_interval_end"].strftime("%Y%m%d"),"%Y%m%d") return time_since_midnight > target_time t=run_after() test()
但我发现Sensor中无法访问上下文变量(context为空字典),而在普通Task中可以正常访问。请问我的操作是否有误?有没有合适的解决方案?(我想到可以在其他Task中读取参数后通过XCom传递给Sensor,但这会让DAG更复杂,似乎不是正确的做法)
解决方案
问题原因
使用@task.sensor装饰器时,默认不会自动注入完整的上下文变量,这是它和普通@task的核心区别之一。普通任务会自动传递上下文,但Sensor装饰器需要显式声明所需的上下文变量,或通过配置开启上下文注入。
可行解决方法
方法1:显式声明所需参数(推荐)
直接在Sensor函数中声明需要的参数(如params、data_interval_end),Airflow会自动将这些上下文变量注入,无需通过**context提取:
@task.sensor( task_id=f"run_after", poke_interval=60 * 5, timeout=60 * 60 * 24 * 3, mode="reschedule" ) def run_after(params, data_interval_end): run_after_time = params.get("time","0h") print(run_after_time) target_time = parse_time(run_after_time) time_since_midnight = datetime.now() - datetime.strptime(data_interval_end.strftime("%Y%m%d"),"%Y%m%d") return time_since_midnight > target_time
这种方式符合Airflow 2.x的设计规范,代码更简洁易读。
方法2:开启上下文注入
在@task.sensor装饰器中添加provide_context=True参数,让**context能获取完整的上下文变量:
@task.sensor( task_id=f"run_after", poke_interval=60 * 5, timeout=60 * 60 * 24 * 3, mode="reschedule", provide_context=True ) def run_after(**context): run_after_time = context["params"].get("time","0h") print(run_after_time) target_time = parse_time(run_after_time) time_since_midnight = datetime.now() - datetime.strptime(context["data_interval_end"].strftime("%Y%m%d"),"%Y%m%d") return time_since_midnight > target_time
注意:该参数在Airflow 2.x中已被标记为deprecated,官方推荐使用方法1的显式参数声明方式。
方法3:全局变量替代(仅适用于非动态场景)
如果参数不需要每次触发DAG时修改,可将其存储在Airflow的Variable中,直接在Sensor中读取:
from airflow.models import Variable # 先在Airflow UI中设置变量time_config的值为"8h" @task.sensor( task_id=f"run_after", poke_interval=60 * 5, timeout=60 * 60 * 24 * 3, mode="reschedule" ) def run_after(): run_after_time = Variable.get("time_config", default_var="0h") target_time = parse_time(run_after_time) time_since_midnight = datetime.now() - datetime.strptime(datetime.now().strftime("%Y%m%d"),"%Y%m%d") return time_since_midnight > target_time
此方法不适合需要动态修改参数的场景,仅适用于全局固定配置。
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

