【问题标题】:Datastore poor performance with Apache Beam & DataflowApache Beam 和 Dataflow 的数据存储性能不佳
【发布时间】:2018-07-01 09:50:03
【问题描述】:

我在数据存储区写入速度方面遇到了巨大的性能问题。大多数时候它保持在 100 个元素/秒以下。

当使用数据存储客户端 (com.google.cloud:google-cloud-datastore) 在本地计算机上对写入速度进行基准测试并并行运行批量写入时,我能够达到大约 2600 个元素/秒的速度.

我已经使用 Java API 设置了一个简单的 Apache Beam 管道。这是它的图表:

以下是在没有 Datastore 节点的情况下运行时的速度:

这种方式要快得多。这一切都表明 DatastoreV1.Write 是这个管道的瓶颈——从没有写入节点的管道速度和 DatastoreV1.Write 的挂墙时间与其他节点的挂墙时间相比来看。


我尝试解决的方法:

• 增加初始工人的数量(尝试了 1 和 10,没有明显差异)。一段时间后,数据存储将写入次数减少到 1(可能在前 2 个节点完成处理之后)。基于此,我怀疑 DatastoreIO.v1().write() 不会并行运行其工作程序。为什么呢?

• 确保所有内容都在同一位置运行:GCP 项目、数据流管道工作程序和元数据、存储 - 都设置为 us-central。这是建议here

• 尝试使用不同的实体密钥生成策略(根据this post)。目前使用这种方法:Key.Builder keyBuilder = DatastoreHelper.makeKey("someKind", UUID.randomUUID().toString());。我不能完全确定这会生成足够均匀分布的密钥,但我想即使它没有,性能也不应该那么低?


请注意,如果没有解决方法,我无法使用提供的 Apache Beam 和 Google 库:由于依赖问题,我不得不强制 google-api-client 版本为 1.22.0 和 Guava 为 23.0(请参阅例如https://github.com/GoogleCloudPlatform/DataflowJavaSDK/issues/607)。

查看DatastoreV1.Write节点日志:

它大约每 5 秒推送 500 个实体,这不是很快。

总体而言,DatastoreIO.v1().write() 速度看起来很慢,并且它的工作人员没有并行运行。知道如何解决这个问题或可能是什么原因吗?

【问题讨论】:

    标签: java google-cloud-platform google-cloud-datastore google-cloud-dataflow apache-beam


    【解决方案1】:

    我不应该不回答这个问题。

    在联系 GCP 支持后,有人向我提供了一个建议,即原因可能是 TextIO.Read 节点从压缩(gzipped)文件中读取。显然这是一个不可并行化的操作。事实上,在切换到源文件的未压缩文件后,性能得到了改善。

    建议的另一个解决方案是在从源代码读取后对管道进行手动重新分区。这意味着向管道中的项目添加任意键,按任意键分组,然后删除任意键。它也可以。这种方法归结为这段代码:

    管道代码:

    pipeline.apply(TextIO.read().from("foo").withCompression(Compression.GZIP)  
            .apply(ParDo.of(new PipelineRepartitioner.AddArbitraryKey<>()))
            .apply(GroupByKey.create())
            .apply(ParDo.of(new PipelineRepartitioner.RemoveArbitraryKey<>()))
            /* further transforms */ 
    

    助手类:

    public class PipelineRepartitioner<T> {
        public static class AddArbitraryKey<T> extends DoFn<T, KV<Integer, T>> {
            @ProcessElement
            public void processElement(ProcessContext c) {
                c.output(KV.of(ThreadLocalRandom.current().nextInt(), c.element()));
            }
        }
    
        public static class RemoveArbitraryKey<T> extends DoFn<KV<Integer, Iterable<T>>, T> {
            @ProcessElement
            public void processElement(ProcessContext c) {
                for (T s : c.element().getValue()) {
                    c.output(s);
                }
            }
        }
    }
    

    我在 Apache Beam Jira 上看到了与该问题相关的票证,因此将来可能会解决此问题。

    【讨论】:

      猜你喜欢
      • 2021-04-29
      • 1970-01-01
      • 1970-01-01
      • 2022-11-26
      • 2020-03-22
      • 1970-01-01
      • 1970-01-01
      • 2022-11-02
      相关资源
      最近更新 更多