【问题标题】:BigqueryIO file loads: only use additional shard if requiredBigqueryIO 文件加载:仅在需要时使用附加分片
【发布时间】:2020-06-09 13:47:29
【问题描述】:

我有一个数据流作业,它从 pubsub 读取数据,将 PubsubMessage 转换为 TableRow,并使用 FILE_LOAD-方法(每 10 分钟,1 个分片)将此行写入 BQ。这项工作有时会引发ByteString would be too long-异常。当它将行连接到 Google Cloud Storage (GCS) 临时文件时,应引发此异常,因为您无法附加到 GCS 文件。如果我理解正确,可以让这个异常发生,因为稍后将使用“大”临时文件加载到 BQ,并且附加将发生在应该成功的新文件上。但是,当我接近项目的每日负载作业配额时,我希望在不增加负载作业数量的情况下防止发生此错误。

我可以:

  • 将分片数增加到 2 个?或者这是否会导致写入器始终使用 2 个分片,即使它只需要写入少量行?
  • 使用setMaxFileSize() 和分片数?或者作者是否仍会使用 2 个碎片,即使它实际上并没有?

提前致谢!

【问题讨论】:

    标签: python google-cloud-dataflow apache-beam


    【解决方案1】:

    将分片数量设置为 2 将始终使用 2 个分片。

    但是,我不认为“ByteString 将太长”错误来自 GCS。当 Dataflow 中捆绑包的总输出大小过大 (>2GB) 时,通常会发生该错误,当 DoFn 的输出远大于其输入时,就会发生这种情况。

    解决此问题的一种方法是使用 GroupByKey 拆分来自 Pubsub 的捆绑包。您可以使用输入的哈希值或随机数作为键,并将触发器设置为 AfterPane.elementCountAtLeast(1) 以允许元素一到达就输出。

    【讨论】:

    • 感谢您的回答!我实现了 GroupByKey,但我不确定如何将它与AfterPane.elementCountAtLeast(1) 结合起来。您能否提供一个示例或指出一些示例或页面来解释如何执行此操作?
    猜你喜欢
    • 1970-01-01
    • 2014-10-31
    • 2013-11-30
    • 1970-01-01
    • 1970-01-01
    • 2020-09-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多