【问题标题】:How to set up job dependencies in google bigquery?如何在 google bigquery 中设置作业依赖项?
【发布时间】:2020-01-31 09:40:38
【问题描述】:

我有一些工作,比如说一个是从谷歌云存储桶加载一个文本文件到 bigquery 表,另一个是预定查询,通过一些转换将数据从一个表复制到另一个表,我想要第二个工作取决于第一个的成功,如果可能的话,我们如何在 bigquery 中实现这一点?

非常感谢。

最好的问候,

【问题讨论】:

  • 您现在需要编写脚本。当 BigQuery 作业完成时,有一个高投票的功能请求来获取事件,因此您可以自动化链条。
  • 你好,你能提供一个如何编写脚本的例子吗?用什么语言?谢谢。

标签: google-bigquery


【解决方案1】:

现在,开发人员需要整合操作链。 可以使用 Cloud Functions(支持、Node.js、Go、Python)或通过 Cloud Run 容器(支持 gcloud API、任何编程语言)来完成。

基本上你需要

  1. 发布作业
  2. 获取工作 ID
  3. 职位 ID 投票
  4. 作业完成触发其他步骤

如果使用云函数

  1. 将文件放入专用的 GCS 存储桶中
  2. 设置一个 GCF 来监控该存储桶,当上传新文件时,它将执行导入 GCS 的函数 - 等待操作结束
  3. 在 GCF 结束时,您可以触发其他功能以进行下一步

Cloud Functions 的另一个用例:

A:触发器启动 GCF
B:函数执行查询(将数据复制到另一个表)
C: 得到一个工作 id - 稍微延迟触发另一个函数

I:一个函数得到一个jobid
J:工作的投票准备好了吗?
K:如果没有准备好,稍微延迟一下再开火
L:如果准备好触发下一步 - 可以是专用函数或参数化函数

【讨论】:

    【解决方案2】:

    可以使用云功能 (CF) 或调度程序 (airflow) 来解决您的场景。第一种方法是事件驱动,让您的数据立即得到处理。使用调度程序,预计数据可用性延迟。

    如前所述,一旦您提交 BigQuery 作业,您将获得作业 ID,需要对其进行检查,直至完成。然后根据您可以分别处理成功或失败发布操作的状态。

    如果您要开发 CF,请注意存在某些限制,例如执行时间(最长 9 分钟),如果 BigQuery 作业需要超过 9 分钟才能完成,您必须解决这些限制。 CF 的另一个挑战是幂等性,确保如果同一个数据文件事件不止一次出现,处理不应导致数据重复。

    或者,您可以考虑使用一些事件驱动的无服务器开源项目,例如 BqTail - Google Cloud Storage BigQuery Loader with post-load transformation。

    这是 bqtail 规则的示例。

    rule.yaml

    When:
      Prefix: "/mypath/mysubpath"
      Suffix: ".json"
    Async: true
    Batch:
      Window:
        DurationInSec: 85
    Dest:
      Table: bqtail.transactions
      Transient:
        Dataset: temp
        Alias: t
      Transform:
        charge: (CASE WHEN type_id = 1 THEN t.payment + f.value WHEN type_id = 2 THEN t.payment * (1 + f.value) END)
      SideInputs:
        - Table: bqtail.fees
          Alias: f
          'On': t.fee_id = f.id
    OnSuccess:
      - Action: query
        Request:
          SQL: SELECT
            DATE(timestamp) AS date,
            sku_id,
            supply_entity_id,
            MAX($EventID) AS batch_id,
            SUM( payment) payment,
            SUM((CASE WHEN type_id = 1 THEN t.payment + f.value WHEN type_id = 2 THEN t.payment * (1 + f.value) END)) charge,
            SUM(COALESCE(qty, 1.0)) AS qty
            FROM $TempTable t
            LEFT JOIN bqtail.fees f ON f.id = t.fee_id
            GROUP BY 1, 2, 3
          Dest: bqtail.supply_performance
          Append: true
        OnFailure:
          - Action: notify
            Request:
              Channels:
                - "#e2e"
              Title: Failed to aggregate data to supply_performance
              Message: "$Error"
        OnSuccess:
          - Action: query
            Request:
              SQL: SELECT CURRENT_TIMESTAMP() AS timestamp, $EventID AS job_id
              Dest: bqtail.supply_performance_batches
              Append: true
          - Action: delete
    

    【讨论】:

      【解决方案3】:

      您想使用编排工具,尤其是当您想将此任务设置为重复作业时。 我们使用Google Cloud Composer,这是一个基于Airflow 的托管服务,用于工作流编排,效果很好。它具有自动重试、监控、警报等功能。

      您可能想尝试一下。

      【讨论】:

        【解决方案4】:

        基本上,您可以使用 Cloud Logging 了解 GCP 中几乎所有类型的操作。

        BigQuery 也不例外。查询作业完成后,您可以在日志查看器中找到对应的日志。

        下一个问题是如何锚定您想要的确切查询,实现此目的的一种方法是使用带标签的查询(意味着将标签附加到您的查询)[1]。

        例如,您可以使用下面的bq 命令发出带有foo:bar 标签的查询

        bq query \
        --nouse_legacy_sql \
        --label foo:bar \
        'SELECT COUNT(*) FROM `bigquery-public-data`.samples.shakespeare'
        

        然后,当您进入日志查看器并发出以下日志过滤器时,您将找到上述查询生成的确切日志。

        resource.type="bigquery_resource"
        protoPayload.serviceData.jobCompletedEvent.job.jobConfiguration.labels.foo="bar"
        

        下一个问题是如何根据此日志为下一个工作负载发出事件。然后,Cloud Pub/Sub 开始发挥作用。

        根据日志模式发布事件的两种方法是:

        1. 日志路由器:将 Pub/Sub 主题设置为目标 [1]
        2. 基于日志的指标:创建通知渠道为 Pub/Sub [2] 的警报策略

        因此,下一个工作负载可以订阅 Pub/Sub 主题,并在上​​一个查询完成时触发。

        希望对你有帮助~

        [1]https://cloud.google.com/bigquery/docs/reference/rest/v2/Job#jobconfiguration
        [2]https://cloud.google.com/logging/docs/routing/overview
        [3]https://cloud.google.com/logging/docs/logs-based-metrics

        【讨论】:

          猜你喜欢
          • 2022-09-23
          • 2015-11-04
          • 2020-03-23
          • 1970-01-01
          • 1970-01-01
          • 2014-06-09
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多