【发布时间】:2019-06-04 11:27:39
【问题描述】:
我正在使用 Cloud Composer 为到达 GCS 并进入 BigQuery 的文件编排 ETL。我有一个云函数,当文件到达时触发 dag,云函数将文件名/位置传递给 DAG。在我的 DAG 中,我有 2 个任务:
1) 使用DataflowPythonOperator 运行数据流作业,从 GCS 中的文本中读取数据并将其转换并输入到 BQ,然后 2) 移动文件到失败/成功存储桶,具体取决于作业是失败还是成功。
每个文件都有一个文件 ID,它是 bigquery 表中的一列。有时一个文件会被编辑一次或两次(这不是经常发生的流式传输),我希望能够首先删除该文件的现有记录。
我查看了其他气流运算符,但在运行数据流作业之前希望在我的 DAG 中有 2 个任务:
- 根据文件名获取文件 ID(现在我有一个 bigquery 表映射文件名 -> 文件 ID,但我也可以只引入一个用作地图的 json,我想这是否更容易)
- 如果文件 ID 已经存在于 bigquery 表(从数据流作业输出转换后的数据的表)中,请将其删除,然后运行数据流作业,这样我就有了最新信息。我知道一个选择是只添加一个时间戳并且只使用最新的记录,但是因为每个文件可能有 100 万条记录,这不像我每天删除 100 个文件(可能是 1-2 个顶部)这似乎是混乱和混乱。
在数据流作业之后,理想情况下,在将文件移动到成功/失败文件夹之前,我想附加一些“记录”表,说明此时输入了该游戏。这将是我查看发生的所有插入的方式。 我试图寻找不同的方法来做到这一点,我是云作曲家的新手,所以在经过 10 多个小时的研究后,我不清楚这将如何工作,否则我会发布代码以供输入。
谢谢,我非常感谢大家的帮助,如果这不是你想要的那么清楚,我深表歉意,关于气流的文档非常强大,但考虑到云作曲家和 bigquery 相对较新,很难彻底了解如何做一些 GCP 特定的任务。
【问题讨论】:
标签: python google-cloud-platform google-bigquery airflow google-cloud-composer