【问题标题】:pyspark with spark 2.4 on EMR SparkException: Cannot broadcast the table that is larger than 8GB在 EMR SparkException 上使用 spark 2.4 的 pyspark:无法广播大于 8GB 的​​表
【发布时间】:2020-03-15 08:30:19
【问题描述】:

我检查了与此错误相关的其他帖子,但没有发现任何工作。

我正在尝试做的事情:

df = spark.sql("""
SELECT DISTINCT
  action.AccountId
  ...
  ,to_date(date) as Date
FROM sc_raw_report LEFT JOIN adwords_accounts ON action.AccountId=sc_raw_report.customer_id
WHERE date >= to_date(concat_ws('-',2018,1,1))
GROUP BY action.AccountId
  ,Account_Name
  ...
  ,to_date(date)

  ,substring(timestamp,12,2)
""")

df.show(5, False)

然后是 saveAsTable.. 尽管如此它返回一个错误:

py4j.protocol.Py4JJavaError:调用 o119.showString 时出错。 :org.apache.spark.SparkException:在 awaitResult 中抛出异常: [...] 原因:org.apache.spark.SparkException:无法广播大于 8GB 的​​表:13 GB

我已经尝试过: 'spark.sql.autoBroadcastJoinThreshold': '-1'

但它什么也没做。

adwords_account 表非常小,在 sc_raw_report 上打印 df.count() 返回:2022197

emr-5.28.0 火花 2.4.4 我的集群核心:15 个 r4.4xlarge(16 个 vCore,122 GiB 内存,仅 EBS 存储) main:r5a.4xlarge(16 vCore,128 GiB 内存,仅 EBS 存储)

使用 spark-submit --deploy-mode 集群的配置:

--conf spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem --conf fs.s3a.attempts.maximum=30 --conf spark.sql.crossJoin.enabled= true --executor-cores 5 --num-executors 5 --conf spark.dynamicAllocation.enabled=false --conf spark.executor.memoryOverhead=3g --driver-memory 22g --executor-memory 22g --conf spark。 executor.instances=49 --conf spark.default.parallelism=490 --conf spark.driver.maxResultSize=0 --conf spark.sql.broadcastTimeout=3600

有人知道我可以在这里做什么吗?

编辑:附加信息: 升级到 16 个实例或 r4.8xlarge(32CPU,244RAM)也无济于事。

带有 step 的图形,然后它在抛出广播错误之前空闲

执行者在崩溃前不久报告:

配置:

spark.serializer.objectStreamReset  100
spark.sql.autoBroadcastJoinThreshold    -1
spark.executor.memoryOverhead   3g
spark.driver.maxResultSize  0
spark.shuffle.service.enabled   true
spark.rdd.compress  True
spark.stage.attempt.ignoreOnDecommissionFetchFailure    true
spark.sql.crossJoin.enabled true
hive.metastore.client.factory.class com.amazonaws.glue.catalog.metastore.AWSGlueDataCatalogHiveClientFactory
spark.scheduler.mode    FIFO
spark.driver.memory 22g
spark.executor.instances    5
spark.default.parallelism   490
spark.resourceManager.cleanupExpiredHost    true
spark.executor.id   driver
spark.driver.extraJavaOptions   -Dcom.amazonaws.services.s3.enableV4=true
spark.hadoop.fs.s3.getObject.initialSocketTimeoutMilliseconds   2000
spark.submit.deployMode cluster
spark.sql.broadcastTimeout  3600
spark.master    yarn
spark.sql.parquet.output.committer.class    com.amazon.emr.committer.EmrOptimizedSparkSqlParquetOutputCommitter
spark.ui.filters    org.apache.hadoop.yarn.server.webproxy.amfilter.AmIpFilter
spark.blacklist.decommissioning.timeout 1h
spark.sql.hive.metastore.sharedPrefixes com.amazonaws.services.dynamodbv2
spark.executor.memory   22g
spark.dynamicAllocation.enabled false
spark.sql.catalogImplementation hive
spark.executor.cores    5
spark.decommissioning.timeout.threshold 20
spark.hadoop.mapreduce.fileoutputcommitter.cleanup-failures.ignored.emr_internal_use_only.EmrFileSystem true
spark.hadoop.yarn.timeline-service.enabled  false
spark.yarn.executor.memoryOverheadFactor    0.1875

【问题讨论】:

  • 目前,Spark 中的广播变量大小应小于 8GB 是硬性限制。请指定两个数据集的大小。您可以在加入连接列之前重新分区数据集并尝试吗?
  • 数据按 customer_id / date 划分,因此很难确切知道数据集有多大。但是对于 adwords_accounts 它应该是 ~ 5MB 并且在 sc_raw_report 上打印 df.count() 返回 2022197

标签: apache-spark pyspark pyspark-sql


【解决方案1】:

ShuffleMapStage 之后,shuffle 块的一部分需要在driver 处为broadcasted。 请确保 Driver(在您的情况下是 YARN 中的 AM)有足够的内存/开销。

你能发布sc运行时配置吗?

【讨论】:

  • 我正在使用 AWS 的 "maximizeResourceAllocation":"true",我应该关注驱动程序的哪些特性?
  • maximizeResourceAllocation — 这解释了问题(据我猜测,不看日志)aws 所做的是通过master instance’s 允许的内存设置它 - spark.driver.memory 集将其设置为更高的数字并关闭最大资源分配
  • 哦,哇,现在我手动调整它似乎确实有效。这次有一个节目通过了,我正在等待确认步骤才能批准。否则我会更新指标/配置
  • 不幸的是,当我恢复到原始数据量时它失败并出现同样的错误......我正在更新我的问题
  • 试试这个创可贴看看——在 hdfs 中写入左连接的结果,然后在一个单独的应用程序(另一个 spark-submit)中做你的小组。 Lmk 进展如何。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-03-04
  • 2020-03-27
  • 2022-08-09
  • 2020-02-12
  • 1970-01-01
  • 2018-03-07
  • 1970-01-01
相关资源
最近更新 更多