【问题标题】:Inserting into BigQuery via load jobs (not streaming)通过加载作业(不是流式传输)插入 BigQuery
【发布时间】:2019-07-05 00:07:28
【问题描述】:

我希望使用 Dataflow 将数据加载到使用 BQ load jobs 的 BigQuery 表中 - 不是流式传输(对于我们的用例而言,流式传输成本太高)。我看到 Dataflow SDK 内置了对通过 BQ 流插入数据的支持,但我无法在 Dataflow SDK 中找到任何支持开箱即用的加载作业的东西。

一些问题:

1) Dataflow SDK 是否支持 BigQuery 加载作业插入?如果没有,有计划吗?

2) 如果我需要自己动手,有什么好的方法?

如果我必须自己动手,使用 Google Cloud Storage 执行 BQ 加载作业是一个多步骤过程 - 将文件写入 GCS,通过 BQ API 提交加载作业,并(可选)检查状态,直到作业已完成(或失败)。我希望我可以使用现有的 TextIO.write() 功能来写入 GCS,但我不确定如何通过随后调用 BQ API 来提交加载作业(以及可选的后续调用来检查作业的状态,直到它完成)。

另外,我会在流模式下使用 Dataflow,窗口为 60 秒 - 所以我也想每 60 秒执行一次加载作业。

建议?

【问题讨论】:

  • (删除答案,并转换为评论)。在批处理模式下,数据流实际上会写入 GCS,然后启动 BigQuery 批量加载作业以获取数据。然后它应该在管道成功(或失败)后删除 GCS 中的文件,但是有错误(goo.gl/8rY1uk)。在流模式下,它确实会使用流 API。我们在这里谈论什么样的尺寸?流媒体每 200MB 的成本仅为 0.01 美元。也许您可以编写 2 个管道 - 一个写入 GCS(以流模式),另一个以批处理模式收集这些文件并使用 BQ 的批量加载?
  • 有趣,没有意识到 DF 根据管道是批处理还是流式传输而不同地写入 BQ。我在DF SDK中仍然找不到批量BQ加载代码,但也许它没有与SDK一起分发?关于流式传输成本,对于我们的用例来说,价格有点高 - 每天大约 60 亿 1KB 行,对于 BQ 流式插入插入来说大约是 9000 美元/月(这是 1KB 的最小行大小,因为我们的大多数行都是实际上是那个大小的一半)。
  • 拥有 2 条管道的想法怎么样?一个写入 GCS(流),另一个(批量)收集数据并批量加载到 BQ?这对你有用吗?也许其中一位 Google 工程师会跳到这里,并提供更多(可能更好:))建议。
  • 感谢您的建议 - 我现在也在考虑流式传输到 GCS 和批量加载到 BQ,尽管我不喜欢增加的操作开销。我想我会更深入地研究代码,看看我是否可以让 DF 即使在流模式下也能批量加载到 BQ。
  • 我今天碰巧看到了这个帖子——这看起来是一个很棒的功能请求,我将在内部提出。延迟和成本之间存在明显的紧张关系——你能告诉我们什么类型的延迟是可以接受的吗?例如,流式传输到 GCS,然后每 24 小时运行一次批量导入作业就可以了……

标签: google-bigquery google-cloud-dataflow


【解决方案1】:

我不确定您使用的是哪个版本的 Apache Beam,但现在可以使用流管道使用微批处理策略。如果你决定一种或另一种方式,你可以使用这样的东西:

.apply("Saving in batches", BigQueryIO.writeTableRows()
                    .to(destinationTable(options))
                    .withMethod(Method.FILE_LOADS)
                    .withJsonSchema(myTableSchema)
                    .withCreateDisposition(CreateDisposition.CREATE_IF_NEEDED)
                    .withWriteDisposition(WriteDisposition.WRITE_APPEND)
                    .withExtendedErrorInfo()
                    .withTriggeringFrequency(Duration.standardMinutes(2))
                    .withNumFileShards(1);
                    .optimizedWrites());

注意事项

  1. 有 2 种不同的方法:FILE_LOADSSTREAMING_INSERT,如果使用第一种方法,则需要包含 withTriggeringFrequencywithNumFileShards。对于第一个,根据我的经验,最好使用分钟,数量将取决于吞吐量数据量。如果您收到很多,请尽量保持较小,当您增加太多时,我会看到“卡住的错误”。分片可能主要影响您的 GCS 计费,如果您添加太多分片,它将每 x 分钟为每个表创建更多文件。
  2. 如果您的输入数据大小不是很大,则流式插入可以很好地工作,并且成本应该不是什么大问题。在这种情况下,您可以使用STREAMING_INSERT 方法并删除withTriggeringFrequencywithNumFileShards。此外,您可以像 InsertRetryPolicy.retryTransientErrors() 一样添加 withFailedInsertRetryPolicy,这样就不会丢失任何行(请记住,STREAM_INSERTS 不能保证幂等性,因此可能会出现重复)
  3. 您可以在 BigQuery 中检查您的作业并确认一切正常!当您尝试定义触发频率和分片时,请记住使用 BigQuery 的作业策略(我认为是每个表 1000 个作业)。

注意:你可以随时阅读这篇关于高效聚合管道的文章https://cloud.google.com/blog/products/data-analytics/how-to-efficiently-process-both-real-time-and-aggregate-data-with-dataflow

【讨论】:

    【解决方案2】:

    BigQueryIO.write() 在输入 PCollection 有界时始终使用 BigQuery 加载作业。如果您希望它在不受限制的情况下也使用它们,请指定 .withMethod(FILE_LOADS).withTriggeringFrequency(...)

    【讨论】:

      猜你喜欢
      • 2023-03-27
      • 1970-01-01
      • 2018-05-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-09-13
      • 2020-11-27
      相关资源
      最近更新 更多