【问题标题】:apache beam pipeline ingesting "Big" input file (more than 1GB) doesn't create any output file摄取“大”输入文件(超过 1GB)的 apache 光束管道不会创建任何输出文件
【发布时间】: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


    【解决方案1】:

    存在一些问题或可能性。我会按优先顺序回答。

    1. 直接运行器用于测试非常小的数据。它的设计是为了最大程度地保证质量,而性能并不是最重要的。例如:
    • 它会随机打乱数据,以确保您不依赖生产中不存在的排序
    • 它会在每一步之后对数据进行序列化和反序列化,以确保数据能够正确传输(生产运行人员会尽可能避免序列化)
    • 它会检查您是否以禁止的方式变异了元素,这会导致您在生产中丢失数据

    你描述的数据不是很大,正常情况下DirectRunner最终可以处理的。

    1. 您已指定numShards(1),它明确消除了所有并行性。它将导致所有数据在单个线程中组合和处理,因此即使在 DirectRunner 上也会比它可能的要慢。一般来说,您会希望避免人为地限制并行性。

    2. 如果有任何内存不足错误或其他错误阻止处理,您应该会看到很多消息。否则,查看分析和 CPU 利用率以确定处理是否处于活动状态会很有帮助。

    【讨论】:

    • 感谢您的反馈和花时间撰写回复,我在原始帖子上添加了一个编辑部分以反映您的反馈。
    【解决方案2】:

    上面的 Kenn Knowles 已经间接回答了这个问题。直接运行器不用于性能测试,仅用于原型设计和数据的良好形成。它没有任何按分区分割和划分工作的机制,并处理内存中的每个数据集。性能测试应使用其他运行器(如 Flink Runner)进行,它们将提供数据拆分和处理高 IO 瓶颈所需的基础设施类型。

    更新:除了这个问题所解决的问题之外,这里还有一个相关问题:How to deal with (Apache Beam) high IO bottlenecks?

    而这里的问题围绕着弄清楚直接跑步者是否可以处理巨大的数据集(我们已经在这里确定这是不可能的);上面提供的链接指向天气生产运行器(如 flink/spark/cloud 数据流)的讨论,可以在本地处理大量数据集(简短的回答是肯定的,但请在链接上检查自己以进行更深入的讨论) .

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-23
      • 1970-01-01
      • 1970-01-01
      • 2019-01-24
      • 1970-01-01
      • 1970-01-01
      • 2018-12-13
      • 1970-01-01
      相关资源
      最近更新 更多