【问题标题】:how to efficiently make airflow dag definitions database-driven如何有效地使气流 dag 定义数据库驱动
【发布时间】:2019-05-02 19:56:23
【问题描述】:

背景

我有一些从第 3 方 api 中提取数据的 dag。

我们需要提取的帐户会随着时间而改变。要确定要提取哪些帐户,取决于我们可能需要查询数据库或发出 HTTP 请求的过程。

在气流之前,我们将在 python 脚本的开头获取帐户列表。然后我们将遍历帐户列表并将每个帐户拉到文件或我们需要做的任何事情。

但是现在,使用气流,在帐户级别定义任务并让气流处理重试功能和日期范围以及并行执行等是有意义的。

因此,我的 dag 可能看起来像这样:

问题

由于每个帐户都是一个任务,因此每次 dag 解析都需要访问帐户列表。但是由于 dag 文件被频繁解析,因此您不必整天都想查询数据库或等待 REST 调用来处理来自每台机器的每个 dag 解析。这可能是资源密集型的,并且可能会花钱。

问题

有没有一种好方法可以将此类配置信息缓存在本地文件中,最好具有指定的生存时间?

想法

我想过几种不同的方法:

  1. 写入 csv 或 pickle 文件并使用 mtime 过期。
    • 对此的担忧是,如果两个进程同时尝试使文件过期,我可能会发生冲突。我不知道这种可能性有多大,也不知道后果是什么,但可能没什么可怕的。
  2. 为所有此类进程创建一个通用的 sqlite 数据库。应该在第一次访问变量时自动创建。每个配置变量在表中都有一行。使用 last_modified_datetime 列来告知何时过期。
    • 需要更精细的代码和依赖项。
  3. 使用气流变量
    • 这样做的好处是它使用现有的数据库,因此每次查询没有 $ 和合理的网络延迟,但它仍然需要网络往返。
    • 具有在多节点设置中的所有节点相同的优势。
    • 确定何时到期可能会有问题,因此可能会创建配置管理器 dag 以定期更新配置变量。
    • 但这会增加部署和开发过程的复杂性——需要填充变量才能正确定义 DAG——所有开发人员也需要在本地进行管理,而不是更多地创建读取缓存方法。
  4. 子标签?
    • 从未使用过它们,但我怀疑它们可以在这里使用。但无论如何,社区似乎都不鼓励使用它们......

你处理过这个问题吗?你找到一个好的解决方案了吗?这些似乎都不是很好。

【问题讨论】:

    标签: python airflow


    【解决方案1】:

    Airflow 默认的 DAG 解析间隔非常宽松:5 分钟。但即便如此,对于大多数人来说,这也是相当多的,因此如果您的部署距离新 DAG 的到期时间不太近,那么增加它是非常合理的。

    一般来说,我会说在每次 DAG 解析心跳时发出 REST 请求并没有那么糟糕。此外,现在调度过程与解析过程分离,因此不会影响任务调度的速度。 Airflow 会为您缓存 DAG 定义。

    如果您认为仍然有理由将自己的缓存放在上面,我的建议是在定义服务器上缓存,而不是在 Airflow 端。例如,在 REST 端点上使用缓存标头并在需要时自行处理缓存失效。但这可能是一些过早的优化,所以我的建议是开始时不使用它,并且只有在您衡量需要它的令人信服的证据时才实施它。

    编辑:关于 Webserver 和 Worker

    确实,Webserver 也会触发 DAG Parses,但不确定频率如何。可能遵循 guicorn worker 刷新间隔(默认为 30 秒)。默认情况下,工人也会在每个任务开始时执行此操作,但如果您激活酸洗 DAG,则可以保存。不确定这是否是个好主意,我听说这注定会被弃用。

    您可以尝试做的另一件事是将其缓存在 Airflow 进程本身中,记住发出昂贵请求的函数。 Python 有一个内置的 functools (lru_cache),加上酸洗它可能就足够了,而且比其他选项容易得多。

    【讨论】:

    • 是的,过早的优化我只是在考虑这是否可行……对于 REST,你是对的,但我们的主数据库是雪花,如果我们将它用于 dag defs,那么我们正在提交到一整天都有仓库,这是$$$。我认为解决方案是在气流元存储服务器上为此目的创建一个数据库并使用它。
    • 我明白了。是的,这是一个特殊的场景,即使是 5 分钟的间隔也是不可接受的,因为您的存储解决方案具有特定的定价模型。在这种情况下,使用廉价的数据库作为视图是一个很好的出路。
    • 另外一点是,解析dag时不只是在list dag区间;它们也被网络服务器解析;更糟糕的是,我相信它们会在每个任务实例开始时再次被解析。因此,如果您有一个包含 100 个任务的 dag,则在一次运行中,您的 dag 将被解析 100 次。如果相同的“account_list”用于这样的多个 dag,那么这可能是很多请求。达到请求配额也是一个问题。我正在尝试在本地 sqlite 数据库中缓存来解决这个问题..
    • 你说得对,我忘了考虑 webserver 和 worker 触发 dag 解析。我已经更新了答案并为您添加了另一个选项。
    【解决方案2】:

    我也有同样的情况。

    对多个帐户进行 API 调用。最初创建了一个 python 脚本来迭代列表。

    当我开始使用 Airflow 时,考虑过您打算做什么。尝试了您列出的 2 个替代方案。经过一些实验后,如果 HTTP 调用失败,决定使用简单的 try-except 块在 python 中处理重试逻辑。原因是

    1. 一个需要维护的脚本
    2. 更少的气流对象
    3. 使用一个脚本更容易重新启动。 (在 Airflow 中重新启动失败的作业并非轻而易举(没有双关语))

    最终取决于你,这是我的经验。

    【讨论】:

    • 有趣。对我来说,我认为很明显“气流方式”更好——即建立一个操作员分成任务,例如通过 account_id。因为通常复杂的是重试和追赶行为,我们基本上可以让气流处理它,我们的代码基本上归结为“获取这个帐户/天来归档”。但也许我需要对传统方法或混合方法进行更多考虑。
    猜你喜欢
    • 1970-01-01
    • 2017-03-31
    • 1970-01-01
    • 1970-01-01
    • 2023-03-25
    • 2017-01-01
    • 1970-01-01
    • 2020-04-13
    • 1970-01-01
    相关资源
    最近更新 更多