Apache Airflow DAG执行时时间戳重复计算问题求助
问题
编写Apache Airflow的DAG脚本时遇到异常行为:原本期望两个任务打印相同的时间戳,但实际运行时第二个任务的时间戳是其执行时的时间,仿佛DAG在执行test2任务时被重新解析了一次。示例脚本如下:
import pendulum from airflow.sdk import dag, task myTimeStamp = pendulum.now("Europe/Amsterdam").strftime('%Y%m%d_%H%M_%S') @dag( schedule="0/1 * * * *", catchup=False, start_date=pendulum.datetime(2025, 9, 2, tz="Europe/Amsterdam"), is_paused_upon_creation=False, tags=["ssh"], ) def test_local(myDate: str): @task.bash() def test1(myDate1: str) -> str: return f'echo "Message: hi from 1! {myDate1}"; sleep 5;' @task.bash() def test2(myDate2: str) -> str: return f'echo "Message: hi from 2! {myDate2}"' test1(myDate) >> test2(myDate) test_local(myTimeStamp)
解决方案
问题根源在于Airflow会定期重新解析DAG文件(默认每30秒一次),你定义的myTimeStamp是在DAG解析阶段执行的,每次解析都会重新计算时间戳。当test2任务执行时,若DAG已被重新解析,就会拿到新的时间戳。以下是几种解决方法:
方法1:使用DAG执行时间(execution_date)
Airflow每个DAG运行都有固定的执行时间,同一运行实例下所有任务共享该时间,用它生成时间戳能保证一致性:
import pendulum from airflow.sdk import dag, task @dag( schedule="0/1 * * * *", catchup=False, start_date=pendulum.datetime(2025, 9, 2, tz="Europe/Amsterdam"), is_paused_upon_creation=False, tags=["ssh"], ) def test_local(): @task.bash() def test1(execution_date) -> str: timestamp = execution_date.in_timezone("Europe/Amsterdam").strftime('%Y%m%d_%H%M_%S') return f'echo "Message: hi from 1! {timestamp}"; sleep 5;' @task.bash() def test2(execution_date) -> str: timestamp = execution_date.in_timezone("Europe/Amsterdam").strftime('%Y%m%d_%H%M_%S') return f'echo "Message: hi from 2! {timestamp}"' test1() >> test2() test_local()
方法2:通过上游任务传递时间戳
新增一个任务生成时间戳,再传递给test1和test2,确保两个任务使用同一值:
import pendulum from airflow.sdk import dag, task @dag( schedule="0/1 * * * *", catchup=False, start_date=pendulum.datetime(2025, 9, 2, tz="Europe/Amsterdam"), is_paused_upon_creation=False, tags=["ssh"], ) def test_local(): @task() def generate_timestamp(): return pendulum.now("Europe/Amsterdam").strftime('%Y%m%d_%H%M_%S') @task.bash() def test1(myDate1: str) -> str: return f'echo "Message: hi from 1! {myDate1}"; sleep 5;' @task.bash() def test2(myDate2: str) -> str: return f'echo "Message: hi from 2! {myDate2}"' timestamp = generate_timestamp() test1(timestamp) >> test2(timestamp) test_local()
方法3:避免在DAG解析阶段计算动态值
不要将myTimeStamp定义在DAG外部,所有动态计算逻辑都放在任务内部或DAG运行阶段执行,这样就不会因DAG重复解析导致值变化。
内容的提问来源于stack exchange,提问作者Cesar
相关产品推荐
相关产品推荐

