【问题标题】:apache airflow idempotent DAG implementationapache气流幂等DAG实现
【发布时间】:2022-10-15 15:15:39
【问题描述】:

我正在使用以下内容为 API 查询生成 startend 时间:

startTime = datetime.now(pytz.timezone('US/Eastern')) - timedelta(hours = 1)
endTime = datetime.now(pytz.timezone('US/Eastern'))

这很好用,并为 API 查询生成正确的参数。但我注意到如果任务失败并且如果我尝试再次重新运行任务,它会根据 DAG 执行的运行时使用 startTimeendTime 的新值。

我正在尝试弄清楚如何使它更具幂等性,因此如果任务失败,我可以重新运行它,并且将在原始任务执行中使用相同的 startTimeendTime

我已经阅读了一些关于模板和宏的内容,但我似乎无法让它正常工作。

这是任务代码。我正在使用 KubernetesPodOperator。

ant_get_logs = KubernetesPodOperator(
    env_vars={
        "startTime": startTime.strftime('%Y-%m-%d %H:%M:%S'),
        "endTime": endTime.strftime('%Y-%m-%d %H:%M:%S'),
        "timeZone":'US/Eastern',
        "session":'none',
    },

    volumes=[volume],
    volume_mounts=[volume_mount],

    task_id='ant_get_logs',
    image='test:1.0.0',
    image_pull_policy='Always',
    in_cluster=True,
    namespace=namespace,
    name='kubepod_ant_get_logs',
    random_name_suffix=True,
    labels={'app': 'backend', 'env': 'dev'},
    reattach_on_restart=True,
    is_delete_operator_pod=True,
    get_logs=True,
    log_events_on_failure=True,
)

谢谢

【问题讨论】:

  • 你能分享完整的任务代码吗?是 PythonOperator
  • @ozs,我用任务代码更新了我的帖子。

标签: python jinja2 airflow idempotent


【解决方案1】:

我将定义一个计算日期并将其推送到 XCom 的新任务。 然后在 KubernetesPodOperator 中我会拉它。

@task(multiple_outputs=True)
def get_times():
    startTime = datetime.now(pytz.timezone('US/Eastern')) - timedelta(hours=1)
    endTime = datetime.now(pytz.timezone('US/Eastern'))

    return {
        "start_time": startTime.strftime('%Y-%m-%d %H:%M:%S'),
        "end_time": endTime.strftime('%Y-%m-%d %H:%M:%S')
    }

ant_get_logs = KubernetesPodOperator(
    env_vars={
        "startTime": "{{ ti.xcom_pull(task_ids='get_times', key='start_time') }}",
        "endTime": "{{ ti.xcom_pull(task_ids='get_times', key='end_time') }}",
        "timeZone":'US/Eastern',
        "session":'none',
    },

    volumes=[volume],
    volume_mounts=[volume_mount],

    task_id='ant_get_logs',
    image='test:1.0.0',
    image_pull_policy='Always',
    in_cluster=True,
    namespace=namespace,
    name='kubepod_ant_get_logs',
    random_name_suffix=True,
    labels={'app': 'backend', 'env': 'dev'},
    reattach_on_restart=True,
    is_delete_operator_pod=True,
    get_logs=True,
    log_events_on_failure=True,
)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-07-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多