【问题标题】:Prefect how to avoid rerunning a task完善如何避免重新运行任务
【发布时间】:2020-08-14 19:37:40
【问题描述】:

在 Prefect 中,假设我有一些管道为列表中的每个日期运行 f(date),并将其保存到文件中。这是一个非常常见的 ETL 操作。在气流中,如果我运行一次,它将回填所有历史日期。如果我再次运行它,它会知道该任务已经运行,并且只运行任何出现的新任务(即最近日期)。

据我所知,在 Prefect 中,它每天都会运行整个管道,即使前一天完成了 99% 的任务。在不切换到 Prefect Cloud 的情况下,有哪些解决方案可以解决这个问题?你只是在退出之前做一些事情,比如让每个任务缓存它在redis中完成的事情吗?

【问题讨论】:

    标签: python workflow airflow etl prefect


    【解决方案1】:

    Prefect 有许多一流的缓存处理方式,具体取决于您想要多少控制。对于每个任务,您可以指定是否应缓存结果、应缓存多长时间以及应如何使缓存失效(年龄、任务的不同输入、流参数值等)。

    缓存任务的最简单方法是使用targets,它允许您指定任务具有模板化的副作用(通常是本地或云存储中的文件,但也可以是数据库条目、redis 键、或其他任何东西)。在任务运行之前,它会检查副作用是否存在,如果存在,则跳过运行。

    例如,此任务会将其结果写入本地文件,该文件自动以任务名称和当前日期为模板:

    @task(result=LocalResult(), target="{task_name}-{today}")
    def get_data():
        return [1, 2, 3, 4, 5]
    

    只要存在匹配文件,任务就不会重新运行。因为{today} 是目标名称的一部分,所以它将隐式缓存任务的值一天。您还可以在模板中使用参数(例如回填日期)来复制 Airflow 的行为。

    要获得更多控制权,您可以通过在任何任务上设置 cache_forcache_validatorcache_key 来使用 Prefect 的 full cache mechanism。如果设置,任务将以Cached 状态完成,而不是Success 状态。当与 Prefect Server 或 Prefect Cloud 等适当的编排后端配对时,Cached 状态可以通过相同任务(或具有相同 cache_key 的任何任务)的未来运行来查询。未来的任务将返回 Cached 状态作为它自己的结果。

    【讨论】:

    • 如果在第一次完成之前调用第二次任务(即尚未设置缓存)怎么办?
    猜你喜欢
    • 2020-02-04
    • 1970-01-01
    • 2014-09-28
    • 1970-01-01
    • 1970-01-01
    • 2016-12-20
    • 1970-01-01
    • 1970-01-01
    • 2014-11-10
    相关资源
    最近更新 更多