【问题标题】:Fix the NumShards while writing Streaming data to GCS在将流式数据写入 GCS 时修复 NumShards
【发布时间】:2020-02-07 19:32:02
【问题描述】:

实际上,我正在尝试将流数据转储到 BigTable,以防由于解析或任何其他问题而失败,我正在将该记录转储到 GCS。所以我在这里应用固定窗口,但我关心的一件事是 num 碎片。在将数据写入 GCS 时,如何指定 num 分片以及 num 分片的工作原理。

.apply(Window.<String>into(FixedWindows.of(Duration.standardSeconds(30L))))
               .apply(TextIO.write().to("gs:").withWindowedWrites());

如果超过 num shards 限制,是不是 TextIO 会覆盖现有文件。

【问题讨论】:

    标签: google-cloud-platform apache-beam google-cloud-pubsub apache-beam-io


    【解决方案1】:

    分片数量设置不能最终覆盖文件。在这种情况下,它是要写入存储的文件数(每个窗口)。通过修改此值,您可以尝试将所有窗口写入单个文件,也可以将每个元素写入单个文件。

    分片的数量决定了对存储的并行写入次数。因此,在考虑管道的性能时,此设置非常重要。更多数量的分片,将更容易并行化,但会产生大量文件。分片数量越少,创建的文件就越少,但会限制并行度。

    根据梁documentation

    除非您需要特定数量的输出文件,否则不建议设置此值。

    如果不设置此值,将由使用的跑步者决定。例如,DataflowRunner 使用管道中设置的最大工作人员数来设置分片数。

    【讨论】:

    • 如果我将 numShards 设置为 60 并运行流式作业 100 天,那么您说它只会创建 60 个文件并将数据附加到 60 个文件中 100 天。我希望我的假设是正确的。
    • TextIO 会写入一个新文件,所以如果该文件已经存在,则会被覆盖。 By default 窗口文件在名称中包含窗口、窗格和分片信息,因此如果您有 30 秒的窗口,每个窗口有 60 个分片,每分钟您将有大约 120 个文件。在 100 天内,您将拥有(每分钟 120 个文件)*(每天 1440 分钟)* 100 天。此外,据我所知,没有开箱即用的配置可以追加而不是用 TextIO 覆盖
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-27
    • 2018-12-02
    • 2019-10-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多