【问题标题】:Optimization Spark job - Spark 2.1优化 Spark 作业 - Spark 2.1
【发布时间】:2020-02-18 18:53:05
【问题描述】:

我的 spark 作业目前在 59 分钟内运行。我想优化它,以便我花费更少的时间。我注意到作业的最后一步需要很长时间(55 分钟)(请参阅下面 Spark UI 中的 spark 作业的屏幕截图)。

我需要将一个大数据集与一个较小的数据集连接起来,在这个连接的数据集上应用转换(创建一个新列)。

最后,我应该有一个基于PSP 列的数据集重新分区(请参阅下面代码的 sn-p)。我还在最后执行排序(根据 3 列对每个分区进行排序)。

所有细节(基础设施、配置、代码)都可以在下面找到。

我的代码片段:

    spark.conf.set("spark.sql.shuffle.partitions", 4158)

    val uh = uh_months
      .withColumn("UHDIN", datediff(to_date(unix_timestamp(col("UHDIN_YYYYMMDD"), "yyyyMMdd").cast(TimestampType)),
        to_date(unix_timestamp(col("january"), "yyyy-MM-dd").cast(TimestampType))))
"ddMMMyyyy")).cast(TimestampType)))
      .withColumn("DVA_1", date_format(col("DVA"), "dd/MM/yyyy"))
      .drop("UHDIN_YYYYMMDD")
      .drop("january")
      .drop("DVA")
      .persist()

    val uh_flag_comment = new TransactionType().transform(uh)
    uh.unpersist()

    val uh_joined = uh_flag_comment.join(broadcast(smallDF), "NO_NUM")
      .select(
        uh.col("*"),
        smallDF.col("PSP"),
        smallDF.col("minrel"),
        smallDF.col("Label"),
        smallDF.col("StartDate"))
      .withColumnRenamed("DVA_1", "DVA")

    smallDF.unpersist()

    val uh_to_be_sorted = uh_joined.repartition(4158, col("PSP"))
    val uh_final = uh_joined.sortWithinPartitions(col("NO_NUM"), col("UHDIN"), col("HOURMV"))

    uh_final

已编辑 - 重新分区逻辑

    val sqlContext = spark.sqlContext
    sqlContext.udf.register("randomUDF", (partitionCount: Int) => {
      val r = new scala.util.Random
      r.nextInt(partitionCount)
      // Also tried with r.nextInt(partitionCount) + col("PSP")
    })

    val uh_to_be_sorted = uh_joined
        .withColumn("tmp", callUDF("RandomUDF", lit("4158"))
        .repartition(4158, col("tmp"))
        .drop(col("tmp"))
    val uh_final = uh_to_be_sorted.sortWithinPartitions(col("NO_NUM"), col("UHDIN"), col("HOURMV"))

    uh_final

smallDF 是我广播的一个小数据集 (535MB)。

TransactionType 是一个类,我根据 3 列的值(MMEDDEBCREDNMTGP)向我的uh 数据框添加新的字符串元素列,检查这些元素的值使用正则表达式的列。

我之前遇到过很多问题(作业失败),因为没有找到随机播放块。我发现我正在溢出到磁盘并且有很多 GC 内存问题,所以我将“spark.sql.shuffle.partitions”增加到 4158。

为什么是 4158?

Partition_count = (stage input data) / (target size of your partition)

所以Shuffle partition_count = (shuffle stage input data) / 200 MB = 860000/200=4300

我有16*24 - 6 =378 cores availaible。因此,如果我想一次性运行所有任务,我应该将 4300 除以 378,大约是 11。然后 11*378=4158

Spark 版本:2.1

集群配置:

  • 24 个计算节点(工作者)
  • 每个 16 个 vcore
  • 每个节点 90 GB RAM
  • 6 个内核已被其他进程/作业使用

当前 Spark 配置:

-master: 纱线

-执行器-内存:26G

-executor-cores: 5

-驱动内存:70G

-num-executors: 70

-spark.kryoserializer.buffer.max=512

-spark.driver.cores=5

-spark.driver.maxResultSize=500m

-spark.memory.storageFraction=0.4

-spark.memory.fraction=0.9

-spark.hadoop.fs.permissions.umask-mode=007

作业是如何执行的:

我们使用 IntelliJ 构建一个工件(jar),然后将其发送到服务器。然后执行一个 bash 脚本。这个脚本:

  • 导出一些环境变量(SPARK_HOME、HADOOP_CONF_DIR、PATH 和 SPARK_LOCAL_DIRS)

  • 使用上面 spark 配置中定义的所有参数启动 spark-submit 命令

  • 检索应用程序的纱线日志

Spark 用户界面截图

DAG

【问题讨论】:

  • 平均而言,您的任务执行大约需要 5 分钟,但您有一个异常值,需要 49 分钟。这是数据倾斜的症状。
  • 谢谢@Gelerion 我会调查的

标签: apache-spark optimization apache-spark-sql spark-ui


【解决方案1】:

@阿里

根据摘要指标,我们可以说您的数据有偏差(最长持续时间:49 分钟,最长随机读取大小/记录:2.5 GB/23,947,440,平均而言,它需要大约 4-5 分钟,处理时间少于 200 MB/1.2 MM 行)

现在我们知道问题可能是少数分区中的数据偏斜,我认为我们可以通过更改重新分区逻辑val uh_to_be_sorted = uh_joined.repartition(4158, col("PSP")) 来解决这个问题,方法是选择一些东西(比如其他一些列或向 PSP 添加任何其他列)

关于数据倾斜和修复的链接很少

https://dzone.com/articles/optimize-spark-with-distribute-by-cluster-by

https://datarus.wordpress.com/2015/05/04/fighting-the-skew-in-spark/

希望对你有帮助

【讨论】:

  • 谢谢@Naga 我认为我的数据确实没有在我的分区之间均匀分布。你知道我是否可以使用 Spark 数据帧的特定功能吗?否则,我考虑编写一个自定义分区器,将我的数据均匀地分布在分区之间,每次我最终得到一个太大的PSP 分区时都会更改哈希值。是否有意义 ?你知道如何使用数据框来做到这一点吗?
  • @Ali,你可以试试这个val r = new scala.util.Random val uh_to_be_sorted = uh_joined.withColumn("tmp", col("PSP") + r.nextInt(4158)).repartition(4158, col("tmp")).drop(col("tmp") ;我在这里尝试做的是引入一个随机数(在这里添加你可以做任何你喜欢的事情),然后在新列上重新分区以均匀分布数据(与前一个相比);如果这可行,请告诉我,以便我可以更新我的答案
  • 感谢您的回复。但是,如果您这样做,每条记录的值 r.nextInt(4158) 将是相同的:因此,我仍然会有不平衡的分区。我按照您的建议做了同样的事情,但使用 UDF:r.nextInt(4158) 每次 PSP 值相同时都会有所不同。但是,现在需要更多时间。
  • 我刚刚给你的代码是为了在键中引入随机性,以便数据在分区中均匀分布,但是你用 UDF 得到了它。应用 UDF 后,您是否看到数据均匀分布在各个分区中,还是仍然存在数据倾斜?
  • 不幸的是,在那之后我仍然有数据倾斜
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-03-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多