You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 09:13:11