【问题标题】:"java.lang.OutOfMemoryError: Requested array size exceeds VM limit" during pyspark collect_list() execution“java.lang.OutOfMemoryError:请求的数组大小超过 VM 限制”在 pyspark collect_list() 执行期间
【发布时间】:2017-10-03 22:13:09
【问题描述】:

我有一个包含大约 40 列浮点数的 5000 万行的大型数据集。

出于自定义转换的原因,我尝试使用 pyspark 的 collect_list() 函数收集每列的所有浮点值,使用以下伪代码:

for column in columns:
   set_values(column, df.select(collect_list(column)).first()[0])

对于每一列,它执行collect_list() 函数并将值设置到其他一些内部结构中。

我正在运行上述独立集群,其中有 2 个 8 核和 64 GB RAM 的主机,为每个主机的 1 个执行程序分配最大 30 GB 和 6 个核,我在执行过程中遇到以下异常,我怀疑它必须处理收集的数组的大小。

java.lang.OutOfMemoryError: 请求的数组大小超过 VM 限制

我在spark-defaults.conf 中尝试了多种配置,包括分配更多内存、分区号、并行性,甚至是Java 选项,但仍然没有运气。

所以我的假设是 collect_list() 与较大数据集上的 executors/drivers 资源密切相关,还是与这些无关?

有什么设置可以帮助我消除这个问题,否则我必须使用collect() 功能?

【问题讨论】:

    标签: apache-spark pyspark out-of-memory spark-dataframe pyspark-sql


    【解决方案1】:

    collect_list 在您的情况下并不比仅调用 collect 更好。对于大型数据集来说,两者都是非常糟糕的主意。并且很少有实际应用。

    两者都需要与记录数成正比的内存量,而collect_list 只是增加了 shuffle 的开销。

    换句话说 - 如果您别无选择,并且需要本地结构,请使用 select 和 collect 并增加驱动程序内存。它不会让事情变得更糟:

    df.select(column).rdd.map(lambda x: x[0]).collect()
    

    【讨论】:

    • 在没有这 2 个的情况下收集每列的所有值有什么特别的建议吗?
    猜你喜欢
    • 2015-11-03
    • 2018-10-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多