Airflow的GCS传感器使用ds、execution_date参数不更新如何解决
问题根因
ds、execution_date都是Airflow中绑定到单个DAG运行实例的固定属性,同一个DAG Run无论内部的传感器轮询多少次,这两个宏的取值都保持为该DAG Run触发时的时间,不会随探测时间动态更新,所以你观测到目录时间一直固定不变,属于符合Airflow设计逻辑的正常现象。
解决方案
分两种场景适配:
场景1:适配原始业务需求(每日检测YYYYMMDD格式的当日目录)
你原本的需求是每日检测对应日期的目录,不需要每次探测都变更目录日期,仅需要到次日自动切换检测新的日期目录,用以下方案即可:
- 把DAG调度间隔修改为按日调度,无需每分钟触发
- 日期格式用
ds_nodash宏,直接输出YYYYMMDD格式的当日日期
参考代码:
dir_main_name = "archive/files" args = { 'start_date': datetime(2021, 11, 29), } with DAG(dag_id='GCS_check', catchup=False, schedule_interval="@daily", default_args=args) as dag: sensor_check = GoogleCloudStorageObjectSensor( task_id='gcs_polling', bucket=BUCKET_NAME, # ds_nodash自动输出当前DAG Run对应日期的YYYYMMDD格式 object=f"{dir_main_name}/{{{{ ds_nodash }}}}" )
该方案下,每日生成的新DAG Run会自动使用当日日期作为目录名,传感器在当日内轮询时始终检测当日目录,完全符合你的业务要求。
场景2:适配测试需求(每次探测都使用当前实时时间生成目录)
如果你需要每次传感器探测时都使用当前的实时时间生成目录,直接调用Airflow内置的macros.datetime获取当前时间即可,参考代码:
dir_main_name = "archive/files" args = { 'start_date': datetime(2021, 11, 29), } with DAG(dag_id='GCS_check', catchup=False, schedule_interval="*/1 * * * *", default_args=args) as dag: sensor_check = GoogleCloudStorageObjectSensor( task_id='gcs_polling', bucket=BUCKET_NAME, # 每次探测时实时获取当前UTC时间,按需求格式化即可 object=f"{dir_main_name}/{{{{ macros.datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S+00:00') }}}}" )
注意:f字符串中Jinja模板的大括号需要双重转义(写成{{{{ }}}}),如果不想转义也可以直接拼接字符串:object=dir_main_name + "/{{ macros.datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S+00:00') }}"
内容的提问来源于stack exchange,提问作者Pinheiro
相关产品推荐
相关产品推荐

