【问题标题】:Airflow Persistent Data Storage Across DAGs跨 DAG 的气流持久数据存储
【发布时间】:2018-08-27 17:36:19
【问题描述】:

我有几个 DAG 创建临时 AWS EMR 集群,然后在它们完成运行后终止它们。我想创建一个每天运行的新 DAG,并生成当天创建的每个 EMR 集群的报告以及运行时长,并通过电子邮件将此报告发送给不同的人。

我需要存储 EMR 集群 ID 值,以便我的报告生成器拥有当天每个 EMR 集群 ID 的列表。我想知道是否可以修改 Airflow 变量来存储此信息,例如,我可以有一个 Airflow 变量,其中键是“EMR_CLUSTERS”,值是 JSON 字符串,其中包含我想要记录的所有数据。或者我可以使用已经用于写入新表的 Airflow 元数据库,在那里我可以写入这些信息?

在 Airflow 中存储永久数据有哪些选择?

【问题讨论】:

  • 难道不是也可以创建一个简单的 DAG 来查询 AWS API 并请求在特定时间创建的 EMR 集群 - 或使用特定标签?节省您重新创建元数据存储的时间?

标签: airflow


【解决方案1】:

您提到的任何一个选项都可以:

  1. 气流变量
  2. 元数据数据库

第三个选项是网络存储。 如果您正在运行分布式气流,那么您可能将 DAG 存储在网络存储中,并将其安装到工作程序/调度程序/网络服务器中。在这种情况下,将基于文件的报告放在此存储上(并可能通过电子邮件发送等)将是一个可靠的选择。

You could write a plugin 可以与这 3 个中的任何一个一起使用,并且可以显示何时写入/发送的内容。

变量

Easily read/written,但是每天在 IMO 上覆盖它有点草率。

元数据数据库

使用SQLAlchemy 创建和读取/写入存储此信息的表。 您可以通过以下方式获得气流元数据数据库的会话:

from airflow import settings
session = settings.Session()

网络存储

在这种情况下,只需正常读取/写入文件即可。

【讨论】:

  • 对于变量选项,他们的文档只演示了如何从变量中读取而不是从 Python 代码中写入它们。我知道如何使用像airflow variable --set key value 这样的 CLI 来设置它们,但是是否可以从代码中设置它们(除了通过 BashOperator 发送上一个命令之外)?我正在尝试查找 airflow.models.Variables API 文档,但找不到它。
  • 另外,我想知道将 MetaDB 用于此类任务是否是不好的做法?我觉得可能会导致潜在危险的基础设施与代码混合。例如,DevOps 中的某个人可以为 Airflow 环境创建一个新的 MetaDB,却没有意识到它也用于代码(idk 只是我脑海中的一个例子)。我觉得应该有一个 Airflow 标准来处理这种常见的事情,或者用户可能会创建自己的数据库来存储这样的信息?
  • 我们使用大量气流,并为我们自己的插件使用元数据数据库。我们有一个自定义执行器作为插件的一部分,如果它们不存在,它会创建表。如果有人杀死了整个气流数据库,那么插件数据可能不是最大的问题。我想到的第四个选项是将信息放入 S3(因为您使用的是 EMR,所以您可能可以访问 S3)。
  • 我考虑了 S3 选项,这就是我倾向于的选项。但是根据我对 S3 的理解,没有锁定机制,所以我想知道如果我的两个工作人员同时在 S3 中包含的 JSON 文件中记录他们的 EMR 集群详细信息,是否会出现竞争条件。
  • 我会这样写,以便每个作业都写入某个文件/文件夹,该文件/文件夹中包含 run_id。这将保持它的唯一性并避免锁定。 (stackoverflow.com/questions/43345991/…)
猜你喜欢
  • 2018-06-22
  • 2020-05-07
  • 1970-01-01
  • 2020-04-13
  • 1970-01-01
  • 1970-01-01
  • 2017-04-05
  • 2011-08-24
  • 1970-01-01
相关资源
最近更新 更多