【发布时间】: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