【发布时间】:2021-02-15 07:20:50
【问题描述】:
我在 Azure Databricks DBR 7.3 LTS、spark 3.0.1、scala 2.12 上运行以下代码 在 Standard_E4as_v4(32.0 GB 内存、4 个内核、1 个 DBU)VM 的(20 到 35 个)worker 集群上 以及 Standard_DS5_v2 类型的驱动程序(56.0 GB 内存,16 核,3 DBU)
目标是处理约 5.5 TB 的数据
我面临以下异常:“org.apache.spark.SparkException:作业因阶段故障而中止:1165 个任务(4.0 GiB)的序列化结果的总大小大于 spark.driver.maxResultSize 4.0 GiB” 在处理 57071 中的 1163 后,在 6.1 分钟内处理了 148.4 GiB 的数据
我没有收集或传输数据到驱动程序,分区数据会导致这个问题吗? 如果是这样的话:
- 有更好的分区方法吗?
- 如何解决这个问题?
代码:
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._
import spark.implicits._
val w = Window.partitionBy("productId").orderBy(col("@ingestionTimestamp").cast(TimestampType).desc)
val jsonDF = spark.read.json("/mnt/myfile")
val res = jsonDF
.withColumn("row", row_number.over(w))
.where($"row" === 1)
.drop("row")
res.write.json("/mnt/myfile/spark_output")
然后我尝试只在没有转换的情况下再次加载和写入数据,并面临同样的问题,代码:
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._
import spark.implicits._
val jsonDF = spark.read.json("/mnt/myfile")
jsonDF.write.json("/mnt/myfile/spark_output")
【问题讨论】:
标签: scala apache-spark apache-spark-sql databricks azure-databricks