【发布时间】:2020-05-27 20:16:22
【问题描述】:
问题:
我正在编写一个 Apache Beam 管道来将 Avro 文件转换为 Parquet 文件(使用 Spark runner)。在我开始转换大尺寸 Avro 文件 (15G) 之前,一切正常。
用于读取 Avro 文件以创建 PColletion 的代码:
PCollection<GenericRecord> records =
p.apply(FileIO.match().filepattern(s3BucketUrl + inputFilePattern))
.apply(FileIO.readMatches())
.apply(AvroIO.readFilesGenericRecords(inputSchema));
来自我的入口点 shell 脚本的错误消息是:
b'/app/entrypoint.sh: line 42: 8 Killed java -XX:MaxRAM=${MAX_RAM} -XX:MaxRAMFraction=1 -cp /usr/share/tink-analytics-avro-to-parquet/ avro-to-parquet-deploy-task.jar
假设
经过一番调查,我怀疑上面的 AvroIO 代码试图将整个 Avro 文件加载为一个分区,这会导致 OOM 问题。
我的一个假设是:如果我可以在读取 Avro 文件时指定分区数,我们以 100 个分区为例,那么每个分区将仅包含 150M 数据,这应该可以避免 OOM 问题。
我的问题是:
- 这个假设能引导我走向正确的方向吗?
- 如果是这样,我如何在读取 Avro 文件时指定分区数?
【问题讨论】:
-
嗨 - 您可以预先将 avro 文件写入分区,然后文件模式应该处理读取所有符合该模式的文件。测试您的假设的一种方法是增加机器的内存以查看它是否没有用完 RAM。但理想情况下,您希望并行读取小文件,因此拆分文件是有意义的
-
你有完整的 NPE 堆栈跟踪吗?如果有,可以附上吗?
-
嗨@ClaudiuS。 avro 文件是分区的,因此每个 avro 文件大约为 100M。读取和处理 1 个文件没有问题,当我加载所有文件时会出现问题,这就是为什么我怀疑 AvroIO 尝试一次性加载所有数据而不是按分区加载。
-
@fuyi 你在读一个15GB的大文件吗?然后分区在这里不起作用,因为它是一个执行程序,需要确切的 15GB RAM 或大小接近 8GB 的 RAM 和充足的 EBS 卷
-
嗨@NagarajTantri 感谢您的评论。请看上面的评论,输入的是一堆100M大小的文件。
标签: apache-spark apache-beam parquet spark-avro avroio