【发布时间】:2019-12-17 22:18:22
【问题描述】:
我正在尝试使用 Apache Beam 2.16.0 构建管道以处理大量 XML 文件。平均每 24 小时计数 7000 万,在峰值负载时可高达 50 亿。 文件大小从 ~1 kb 到 200 kb 不等(有时可能更大,例如 30 mb)
文件经过各种转换,最终目的地是 BigQuery 表以供进一步分析。因此,首先我读取 xml 文件,然后反序列化为 POJO(在 Jackson 的帮助下),然后应用所有必需的转换。转换工作非常快,在我的机器上,我每秒可以进行大约 40000 次转换,具体取决于文件大小。
我主要关心的是文件读取速度。我觉得所有的阅读都是通过一个工人完成的,我不明白这怎么能并行。我在 10k 测试文件数据集上进行了测试。
我的本地机器(MacBook pro 2018:ssd、16 gb ram 和 6 核 i7 cpu)上的批处理作业每秒可以解析大约 750 个文件。如果我在 DataFlow 上运行它,使用 n1-standard-4 机器,我只能得到大约 75 个文件/秒。它通常不会扩大规模,但即使扩大规模(有时多达 15 个工作人员),我也只能获得大约 350 个文件/秒。
更有趣的是流媒体工作。它立即从 6-7 个工作人员开始,在 UI 上我可以看到 1200-1500 个元素/秒,但通常它不显示速度,如果我选择页面上的最后一项,它表明它已经处理了 10000 个元素。
批处理作业和流作业的唯一区别是 FileIO 的这个选项:
.continuously(Duration.standardSeconds(10), Watch.Growth.never()))
为什么这会对处理速度产生如此大的影响?
运行参数:
--runner=DataflowRunner
--project=<...>
--inputFilePattern=gs://java/log_entry/*.xml
--workerMachineType=n1-standard-4
--tempLocation=gs://java/temp
--maxNumWorkers=100
运行区域和存储桶区域相同。
管道:
pipeline.apply(
FileIO.match()
.withEmptyMatchTreatment(EmptyMatchTreatment.ALLOW)
.filepattern(options.getInputFilePattern())
.continuously(Duration.standardSeconds(10), Watch.Growth.never()))
.apply("xml to POJO", ParDo.of(new XmlToPojoDoFn()));
xml文件示例:
<LogEntry><EntryId>0</EntryId>
<LogValue>Test</LogValue>
<LogTime>12-12-2019</LogTime>
<LogProperty>1</LogProperty>
<LogProperty>2</LogProperty>
<LogProperty>3</LogProperty>
<LogProperty>4</LogProperty>
<LogProperty>5</LogProperty>
</LogEntry>
现实生活中的文件和项目要复杂得多,有很多嵌套节点和大量的转换规则。
GitHub 上的简化代码:https://github.com/costello-art/dataflow-file-io 它只包含“瓶颈”部分——读取文件并反序列化为 POJO。
如果我可以在我的机器(这是一个强大的工作人员)上处理大约 750 个文件/秒,那么我希望在 Dataflow 中类似的 10 个工作人员上处理大约 7500 个文件/秒。
【问题讨论】:
-
有趣...在批处理模式与流模式下,捆绑包的大小通常更大,但不确定这是否相关。您能否尝试使用它,它已被贬低,但这只是为了检查一些东西:FileIO.match() 和您的 XmlToPojoDoFn() 之间的 Reshuffle.ViaRandomKey()
-
@RezaRokni 我没有发现任何区别。但是在测试期间,我更改了读取操作:FileIO.Match -> FileIO.ReadMatches -> 将文件读取为字节 -> 将字节转换为 POJO。我还在批处理模式下对我的生产代码和更大的数据集(144k 和 1m 文件)进行了测试:在 n1-standard-2 上,我能够在 17 名工作人员的情况下获得大约 1000k 文件/秒。这好多了,但我仍然没有像我的本地机器那样接近 850 个文件/秒。需要做更多的测试。我用这种方法更新了 github 代码。
-
顺便说一句,当您在本地机器上进行测试时,文件在哪里?在本地磁盘上还是在云存储桶中?
-
@RezaRokni 本地 (ssd)。我想知道我是否可以在云上实现类似的读取性能
-
我认为您的文件位于 Cloud Storage 中,它们将被拉到工作人员身上。
标签: java performance file-io google-cloud-dataflow apache-beam-io