【问题标题】:Apache Beam AvroIO read large file OOMApache Beam AvroIO 读取大文件 OOM
【发布时间】: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 问题。

我的问题是:

  1. 这个假设能引导我走向正确的方向吗?
  2. 如果是这样,我如何在读取 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


【解决方案1】:

Spark session 没有设置分区数,而是有一个名为spark.sql.files.maxPartitionBytes的属性,默认设置为128Mb,请参阅参考here

Spark 在将输入文件读入内存时使用此编号对输入文件进行分区。

我测试了一个 50Gb 的 avro 文件,Spark 将其分区为 403 个分区。这种 Avro 到 Parquet 的转换适用于具有 16Gb 内存和 4 个内核的 Spark 集群。

【讨论】:

    猜你喜欢
    • 2018-09-07
    • 2020-08-03
    • 1970-01-01
    • 2018-08-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多