【问题标题】:How to deal with (Apache Beam) high IO bottlenecks?如何处理(Apache Beam)高 IO 瓶颈?
【发布时间】:2021-09-13 07:27:53
【问题描述】:

假设举一个简单的例子,我有一个非常简单的光束管道,它只是从文件中读取数据并将数据转储到输出文件中。现在让我们考虑输入文件很大(几 GB 大小,通常无法在文本编辑器中打开的文件类型)。由于直接运行器实现非常简单(它将整个输入集读取到内存中),它无法读取和输出那些巨大的文件(除非您为 java vm 进程分配了不切实际的大量内存);所以我的问题是:“像 flink/spark/cloud dataflow 这样的生产运行程序”如何处理这个“巨大的数据集问题”? - 假设他们不只是尝试将整个文件/数据集放入内存中?” -.

我希望生产运行程序的实现需要“分批或分批”工作(例如分批读取/处理/输出),以避免尝试在任何特定时间点将大量数据集放入内存中。有人可以分享他们对生产运行人员如何处理这种“海量数据”情况的反馈吗?

概括,请注意这也适用于其他输入/输出机制,例如,如果我的输入是来自巨大数据库表的 PCollection(从广义上讲,行大小和数量都很大),生产的内部实现是否runner 以某种方式将给定的输入 SQL 语句划分为许多内部生成的子语句,每个子语句都采用较小的子集(例如,通过内部生成一个 count(-) 语句,然后是 N 个语句,每个语句都采用 (count(-)/N) 个元素?直接-runner 不会这样做,只会将给定的查询 1:1 传递给数据库),或者我作为开发人员有责任“分批迭代”并划分问题,如果确实如此,那是什么这里的最佳实践,即:有一个或多个管道?如果只有一个,然后以某种方式参数化管道以批量读/写?还是迭代一个简单的管道并在管道外部管理必要的元数据?

提前致谢,任何反馈将不胜感激!


编辑(反映大卫的反馈):

David,您的反馈非常有价值,肯定触及了我感兴趣的点。有一个工作发现阶段来拆分源和读取阶段以同时读取拆分分区绝对是我感兴趣的听到,所以感谢您为我指明正确的方向。如果您不介意,我有几个小的后续问题:

1 - 文章在“通用枚举器-阅读器通信机制”部分指出以下内容:

"SplitEnumerator 和 SourceReader 都是用户实现的类。 实施需要一些沟通的情况并不少见 在这两个组件之间。为了促进此类用例 [....]"

所以我的问题是,由某些用户(即开发人员)触发的“拆分 + 阅读行为”提供了实现(特别是 SplitEnumerator 和 SourceReader),还是我可以在没有任何自定义代码的情况下从开箱即用中受益?.

2 - 可能只是更深入地研究上述问题;如果我有批处理/有界工作负载(假设我正在使用 apache flink),并且我有兴趣处理原始帖子中描述的“巨大文件”,那么管道是否会“开箱即用”工作(做幕后的“工作准备阶段”拆分和并行读取),还是需要开发人员实现一些自定义代码?

提前感谢您提供的所有宝贵反馈!

【问题讨论】:

    标签: apache-spark apache-flink google-cloud-dataflow apache-beam dataflow


    【解决方案1】:

    请注意,如果输入是有界的并且事先已知(即批处理工作负载而不是流式处理),这会更直接。

    在设计时考虑到流式传输的 Flink 中,这是通过将“工作发现”与“阅读”分开来完成的。单个SplitEnumerator 运行一次并枚举要读取的块(拆分/分区),并将它们分配给并行读取器。在批处理情况下,拆分由一系列偏移量定义,而在流式传输情况下,每个拆分的结束偏移量设置为 LONG_MAX。

    这在FLIP-27: Refactor Source Interface中有更详细的描述。

    【讨论】:

    • 大卫,您的反馈非常有价值,而且绝对是我正在研究的方向。我在原始帖子中添加了一个编辑部分以反映您的 cmets,您能检查一下吗?谢谢!
    • 我相信开箱即用的体验是您所希望的,但我自己还没有尝试过。实施 FLIP-27 是一项正在进行的工作,我对当前状态并不完全清楚。然而,旧的实现在这些细节上是相似的,除了批处理/有界情况与流情况分开处理。
    • 再次感谢您的反馈,是的;我从 FLIP-27 得到的是,他们计划重构“工作发现”和“读取”功能,使其不会在 SourceFunction 接口和 DataStream API 之间“混合”;但核心功能应该已经存在。再次感谢您的时间和反馈。我将安装 flink 来运行那些非功能测试。
    【解决方案2】:

    只是为了解决这个问题,这个问题的理由是要知道 apache beam - 当与生产运行器结合时 - (如 flink 或 spark 或谷歌云数据流),是否提供开箱即用的机制 -拆分工作又名读/写操作 - 巨大的文件(或一般的数据源)。上面 David Anderson 提供的评论在暗示 Apache flink 如何处理此工作流方面非常有价值。

    此时,我已经使用基于“beam on flink”的管道实现了使用大文件(用于测试可能的 IO 瓶颈)的解决方案,并且我可以确认,flink 将创建一个执行计划,其中包括拆分源和划分以不会出现内存问题的方式工作。现在,当然可能存在稳定性/“IO性能”受到影响的情况,但至少我可以确认在管道抽象后面执行的工作流在执行任务时使用文件系统,避免将所有数据放入内存和从而避免琐碎的内存错误。结论:是的,“beam on flink”(可能还有 spark 和 dataflow)确实提供了适当的工作准备、工作拆分和文件系统使用,以便以有效的方式使用可用的易失性内存。

    关于数据源的更新: 将 DB 视为数据源,Flink 不会(也不能 - 这不是微不足道的)优化/拆分/分配与 DB 数据源相关的工作,就像它优化的方式一样从文件系统读取。虽然仍然有从数据库读取大量数据(记录)的方法,但实现细节需要由开发人员解决,而不是由框架负责。我发现这篇文章 (https://nl.devoteam.com/expert-view/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam/) 非常有助于解决从 Beam 中的数据库读取大量记录的问题(这篇文章使用云数据流运行器,但我使用了 Flink,它工作得很好),拆分查询并分发处理。

    【讨论】:

      猜你喜欢
      • 2019-04-05
      • 1970-01-01
      • 1970-01-01
      • 2015-06-03
      • 2019-07-14
      • 2016-11-12
      • 1970-01-01
      • 1970-01-01
      • 2020-06-15
      相关资源
      最近更新 更多