【问题标题】:Airflow DAG - how to check BQ first (delete if necessary) and then run dataflow job?Airflow DAG - 如何先检查 BQ(必要时删除)然后运行数据流作业?
【发布时间】:2019-06-04 11:27:39
【问题描述】:

我正在使用 Cloud Composer 为到达 GCS 并进入 BigQuery 的文件编排 ETL。我有一个云函数,当文件到达时触发 dag,云函数将文件名/位置传递给 DAG。在我的 DAG 中,我有 2 个任务:

1) 使用DataflowPythonOperator 运行数据流作业,从 GCS 中的文本中读取数据并将其转换并输入到 BQ,然后 2) 移动文件到失败/成功存储桶,具体取决于作业是失败还是成功。 每个文件都有一个文件 ID,它是 bigquery 表中的一列。有时一个文件会被编辑一次或两次(这不是经常发生的流式传输),我希望能够首先删除该文件的现有记录。

我查看了其他气流运算符,但在运行数据流作业之前希望在我的 DAG 中有 2 个任务:

  1. 根据文件名获取文件 ID(现在我有一个 bigquery 表映射文件名 -> 文件 ID,但我也可以只引入一个用作地图的 json,我想这是否更容易)
  2. 如果文件 ID 已经存在于 bigquery 表(从数据流作业输出转换后的数据的表)中,请将其删除,然后运行数据流作业,这样我就有了最新信息。我知道一个选择是只添加一个时间戳并且只使用最新的记录,但是因为每个文件可能有 100 万条记录,这不像我每天删除 100 个文件(可能是 1-2 个顶部)这似乎是混乱和混乱。

在数据流作业之后,理想情况下,在将文件移动到成功/失败文件夹之前,我想附加一些“记录”表,说明此时输入了该游戏。这将是我查看发生的所有插入的方式。 我试图寻找不同的方法来做到这一点,我是云作曲家的新手,所以在经过 10 多个小时的研究后,我不清楚这将如何工作,否则我会发布代码以供输入。

谢谢,我非常感谢大家的帮助,如果这不是你想要的那么清楚,我深表歉意,关于气流的文档非常强大,但考虑到云作曲家和 bigquery 相对较新,很难彻底了解如何做一些 GCP 特定的任务。

【问题讨论】:

    标签: python google-cloud-platform google-bigquery airflow google-cloud-composer


    【解决方案1】:

    听起来有点复杂。很高兴,几乎所有 GCP 服务都有运营商。另一件事是何时触发 DAG 执行。你想通了吗?您希望每次有新文件进入该 GCS 存储桶时触发 Google Cloud Function 运行。

    1. 触发 DAG

    要触发 DAG,您需要使用依赖于 Object FinalizeMetadata Update 触发器的 Google Cloud 函数来调用它。

    1. 将数据加载到 BigQuery

    如果您的文件已经采用 GCS 格式,并且采用 JSON 或 CSV 格式,那么使用 Dataflow 作业就显得多余了。您可以使用GoogleCloudStorageToBigQueryOperator 将文件加载到BQ。

    1. 跟踪文件 ID

    计算文件 ID 的最佳方法可能是使用来自 Airflow 的 Bash 或 Python 运算符。你能直接从文件名中推导出来吗?

    如果是这样,那么您可以在GoogleCloudStorageObjectSensor 的上游使用 Python 运算符来检查文件是否在成功的目录中。

    如果是,那么您可以使用BigQueryOperator 在 BQ 上运行删除查询。

    之后,您运行 GoogleCloudStorageToBigQueryOperator。

    1. 移动文件

    如果您要将文件从 GCS 移动到 GCS 位置,那么 GoogleCloudStorageToGoogleCloudStorageOperator 应该可以满足您的需要。如果您的 BQ 加载运算符失败,则移动到失败的文件位置,如果成功,则移动到成功的作业位置。

    1. 记录任务日志

    也许您需要跟踪插入的只是将任务信息记录到 GCS。查看how to log task information to GCS

    这有帮助吗?

    【讨论】:

    • 所以我已经完成了云功能步骤和数据流工作(数据流是因为它有 100 万条记录,并且气流没有那么快,除非我错了)。问题是我需要能够在计算文件 ID 后访问它。我可以通过使用文件名并在 bigquery 表中进行简单查找来计算它。不过,在我计算之后,我需要我的模板来知道能够在我的查询中使用它是什么。这有意义吗?
    • Airflow 不会导入文件本身。它会触发文件的 BQ 加载作业 - 这非常快速且免费(与 Dataflow 作业不同)。 - 至于文件ID,我一会儿回复你
    • 你能做到吗 (1) 将数据导入临时表 - (2) 等待 1 小时进行文件更新 (3) 如果文件更新,则删除临时表 (4) 如果文件没有更新,将临时表复制到目标表?
    • 很遗憾,没有,而且通常更新不会在 1 小时内,可能会在 24 小时后,也可能会在 1 个月甚至一年后。
    • 在这种情况下,是的 - 您应该有一个 (1) Python 运算符来计算文件 ID。 (2) 对具有该文件 ID 的所有行运行删除查询的 BigQueryOperator。 (3)。运行将文件 ID 添加到行并插入 t BigQuery 的作业的 DataflowPythonOperator。那怎么样?
    猜你喜欢
    • 2020-01-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-21
    • 1970-01-01
    • 2017-06-07
    相关资源
    最近更新 更多