【问题标题】:Sparklyr configuration results in java.lang.OutOfMemoryErrorSparklyr 配置导致 java.lang.OutOfMemoryError
【发布时间】:2021-11-03 14:32:38
【问题描述】:

我在具有 8 个内核和 64Gb RAM 的本地实例上使用 R 运行 sparklyr。我的工作是 left_join 一个 [50 000 000, 12] 数据帧和一个 [20 000 000, 3] 数据帧,我使用 Spark 运行。

# Load packages
library(tidyverse)
library(sparklyr)


# Initialize configuration with defaults
config <- spark_config()

# Memory
# Set memory allocation for whole local Spark instance
# Sys.setenv("SPARK_MEM" = "50g")

# Set driver and executor memory allocations
# config$spark.driver.memory <- "8g"
# config$spark.driver.maxResultSize <- "8g"


# Connect to local cluster with custom configuration
sc <- spark_connect(master = "local", config = config, spark_home = spark_home_dir())


# Read df1 and df2
df1 <- spark_read_parquet(sc, 
                          path = "/mnt/df1/",
                          memory = FALSE, overwrite = TRUE)
df2 <- spark_read_parquet(sc, 
                          path = "/mnt/df2/",
                          memory = FALSE, overwrite = TRUE)


# Left join
df3 <- df1 %>%
  dplyr::left_join(df2)


# Write or collect
sparklyr::spark_write_parquet(df3, path="/mnt/") # or
df3 <- df3 %>% collect()

无论我如何配置 Spark 配置文件,代码都会失败并显示 java.lang.OutOfMemoryError: Java heap space

Error: org.apache.spark.SparkException: Job aborted due to stage failure: Task 2 in stage 8.0 failed 1 times, most recent failure: Lost task 2.0 in stage 8.0 (TID 96, localhost, executor driver): java.lang.OutOfMemoryError: Java heap space

到目前为止,我已经尝试了不同的组合

Sys.setenv("SPARK_MEM" = "50g")
config["sparklyr.shell.driver-memory"] <- "20G"  
config["sparklyr.shell.num-executors"] <- 8  
config$spark.driver.maxResultSize <- "8g" 
config$spark.executor.memory <- "8g"  
config$spark.memory.fraction <- 0.9 

在 R 脚本或 spark 配置文件中。

类似的问题已经被问到1 2 3 但这些都没有解决我的问题。

【问题讨论】:

    标签: r sparklyr


    【解决方案1】:

    您必须为left_join() 指定一个连接键。

    否则,您将尝试计算笛卡尔积,其大小为 [1 000 000 000 000 000, 15] 并且肯定会溢出内存。

    另外:避免在大型数据集上调用 collect(),因为这会将所有数据移回驱动程序,并且驱动程序很可能会出现 OOM。

    【讨论】:

    • 谢谢皮埃尔。连接发生在两个数据框中具有相同名称的两个变量上。我假设 R 只尝试加入具有相似名称的变量。但是,我再次尝试指示要加入 by = c("var1", "var2") 的列,但没有进一步成功。我知道应该谨慎使用 collect() 函数,但最终的数据帧有 dim[50 000 000, 13] 并且应该适合内存。无论我运行 collect 还是尝试在磁盘上写入,都会发生错误。因此,问题必须与加入功能。
    • 您对默认连接键是正确的:我忘记了 sparklyr 和 SparkR 之间的区别。但是如果键分布是倾斜的(即几个键非常频繁),连接的输出仍然可能非常大。加入后,您可以探索最大的密钥:df_stats = df1 %&gt;% count(var1, var2, name="n1") %&gt;% left_join(df2 %&gt;% count(var1, var2, name="n2")) %&gt;% mutate(n3=n1*coalesce(n2,1)) %&gt;% arrange(desc(n3))。你有很大n3的钥匙吗?
    • 感谢您的意见!我试过了,但由于内存问题导致脚本失败java.lang.OutOfMemoryError: Java heap space还有其他建议吗?
    • 很奇怪。直接在 scala spark-shell 中编写查询时,您是否也有同样的错误?也许您可以尝试拆分此计数查询:s1 = df1 %&gt;% count(var1, var2, name="n1") %&gt;% arrange(desc(n1)) %&gt;% collect() & 与 df2 相同,然后加入。你会看到哪个部分有问题。
    • 这适用于df2,其中所有n1 值都是1。对于df1,这将失败。你会推荐在 RStudio 中运行 spark-shell 吗?如果是的话,你有什么指示吗?
    猜你喜欢
    • 2017-04-16
    • 2012-04-10
    • 1970-01-01
    • 2019-03-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-18
    • 2011-04-20
    相关资源
    最近更新 更多