【发布时间】:2018-09-19 15:47:49
【问题描述】:
目前使用带有 Python 的 Google Dataflow 进行批处理。这很好用,但是,我有兴趣在不必处理 Java 的情况下从我的数据流作业中获得更快的速度。
使用 Go SDK,我实现了一个简单的管道,它从 Google 存储中读取一系列 100-500mb 文件(使用 textio.Read),进行一些聚合并使用结果更新 CloudSQL。正在读取的文件数量可以从几十个到几百个不等。
当我运行管道时,我可以从日志中看到文件正在被串行读取,而不是并行读取,因此作业需要更长的时间。使用 Python SDK 执行的同一过程会触发自动缩放并在几分钟内运行多次读取。
我尝试使用 --num_workers= 指定工作人员的数量,但是,Dataflow 会在几分钟后将作业缩减为一个实例,并且从日志中看,在实例运行期间没有发生并行读取。
如果我删除 textio.Read 并实现自定义 DoFn 以从 GCS 读取,则会发生类似的情况。读取过程仍然是串行运行的。
我知道当前的 Go SDK 是实验性的并且缺少许多功能,但是,我还没有找到关于并行处理限制的直接参考,here。当前版本的 Go SDK 是否支持 Dataflow 上的并行处理?
提前致谢
【问题讨论】:
标签: go google-cloud-platform google-cloud-dataflow apache-beam