【发布时间】:2021-06-25 09:39:48
【问题描述】:
关于计算的数据流模型,我正在做一个 PoC 以使用带有直接运行器(和 java sdk)的 apache Beam 来测试一些概念。我无法创建一个读取“大”csv 文件(大约 1.25GB)并将其转储到输出文件中的管道,而无需像以下代码中那样进行任何特定转换(我主要关心使用此数据流测试 IO 瓶颈/beam 模型,因为这对我来说最重要):
// Example 1 reading and writing to a file
Pipeline pipeline = Pipeline.create();
PCollection<String> output = ipeline
.apply(TextIO.read().from("BIG_CSV_FILE"));
output.apply(
TextIO
.write()
.to("BIG_OUTPUT")
.withSuffix("csv").withNumShards(1));
pipeline.run();
我遇到的问题是只有较小的文件才能工作,但是当使用大文件时,不会生成输出文件(但也没有显示错误/异常,这使得调试更加困难)。
我知道在 apache-beam 项目 (https://beam.apache.org/documentation/runners/direct/) 的 runners 页面上,它在内存考虑点下明确说明:
本地执行受到本地环境中可用内存的限制。强烈建议您使用 数据集足够小以适合本地内存。你可以创建一个小 使用创建转换的内存数据集,或者您可以使用读取 转换为使用小型本地或远程文件。
上面的这表明我有一个内存问题(但遗憾的是没有在控制台上明确说明,所以我只是想知道这里)。我也担心他们的建议,即数据集应该适合内存(为什么不从文件中读取部分而不是将整个文件/数据集放入内存?)
我还想在此对话中添加的第二个考虑因素是(如果这确实是内存问题):直接运行器的实现有多基本?我的意思是,实现一段从大文件中分块读取并输出到新文件(也是分块)的代码并不难,因此内存使用在任何时候都不会成为问题(因为这两个文件都没有完全加载到内存中——只有当前的“块”)。即使“直接运行器”更像是一个测试语义的原型运行器,期望它能够很好地处理大文件是否太过分了? - 考虑到这是一个统一模型,专为处理流式传输而构建,其中窗口大小是任意的,并且在下沉之前需要大量数据累积/聚合,这是一个标准用例。
所以不仅仅是一个问题,我非常感谢您就以下任何一点提供反馈/cmets:您是否注意到使用直接运行器的 IO 限制?我是否忽略了某些方面,或者直接跑步者真的如此幼稚地实施?您是否已验证通过使用像 flink/spark/google cloud dataflow 这样的适当生产运行器,此约束消失了?
我最终会与其他运行程序(例如 flink 或 spark 运行程序)一起进行测试,但直接运行程序(即使它仅用于原型制作)在我正在运行的第一个测试中遇到了问题,这让我感到难以接受on - 考虑到整个数据流的想法是基于在统一的批处理/流模型的保护伞下摄取、处理、分组和分发大量数据。
编辑(反映 Kenn 的反馈): Kenn,感谢您提供的宝贵意见和反馈,它们在向我指出相关文档方面提供了很大帮助。根据您的建议,我通过分析应用程序发现问题确实是与 java 堆相关的问题(以某种方式从未在普通控制台上显示 - 仅在分析器上看到)。即使文件“只有”1.25GB 大小,内部使用量在转储堆之前超过 4GB,这表明直接运行程序不是“按块工作”,而是确实将所有内容加载到内存中(正如他们的文档所说)。
关于你的观点:
1- 我相信序列化和洗牌仍然可以通过“逐块”实现来实现。也许我对直接跑步者应该具备的能力有一个错误的期望,或者我没有完全掌握它的预期范围,现在我将避免在使用直接跑步者时进行非功能类型的测试。
2 - 关于分片。我相信 NumOfShards 在写入阶段控制并行度(和输出文件的数量)(在此之前的处理应该仍然是完全并行的,并且只有在写作时,它才会使用尽可能多的工作人员 - 并生成尽可能多的文件 -明确规定)。相信这一点的两个原因是:首先,CPU 分析器总是显示 8 个忙碌的“直接运行者”——反映我的 PC 拥有的逻辑核心的数量——独立于我设置 1 个分片还是 N 个分片。第二个原因是我从这里的文档中了解到的 (https://beam.apache.org/releases/javadoc/2.0.0/org/apache/beam/sdk/io/WriteFiles.html):
默认情况下,输入 PCollection 中的每个包都将由 一个 FileBasedSink.WriteOperation,所以输出的数量会有所不同 基于跑步者的行为,尽管至少 1 个输出将始终是 产生。 可以控制写入阶段的精确并行度 使用 withNumShards(int),通常用于控制文件的数量 生产或在全球范围内限制连接的工人数量 外部服务。但是,此选项通常会损害性能: 它向管道添加了一个额外的 GroupByKey。
这里有一个有趣的事情是“添加到管道中的额外 GroupByKey”在我的用例中是不受欢迎的(我只希望结果在 1 个文件中,不考虑顺序或分组), 因此,在生成 N 个分片输出文件之后,添加一个额外的“展平”文件步骤可能是一种更好的方法。
3 - 您对分析的建议是正确的,谢谢。
Final Edit直接运行器不用于性能测试,仅用于原型设计和数据的良好形成。它没有任何按分区拆分和划分工作的机制,并且处理内存中的所有内容
【问题讨论】:
标签: java google-cloud-dataflow apache-beam