【问题标题】:How to decrease total timing processing of Spark SQL Execution plan如何减少 Spark SQL 执行计划的总时间处理
【发布时间】:2021-07-04 08:21:08
【问题描述】:

我刚刚开发了一个 Spark SQL 应用程序,在一些算法分析过程中,我意识到执行计划需要大量时间来处理。如何优化 Spark SQL 执行计划的性能?

我查看了我们社区中有关此问题的几个问题/答案,但在我看来,没有任何内容能直截了当地执行此操作。因此,我希望得到一些社区支持来克服我的障碍,并可能在此过程中为进一步的开发人员留下路线图。

以下是有关所做工作的一些详细信息。

我开发了一个 Spark 应用程序,它会定期从 kafka 获取事件,对其进行处理并将输出再次发送回 kafka。简而言之,Spark 算法过滤/丰富信息并针对每个事件执行繁重而复杂的窗口滞后函数。

Spark 算法循环运行,因此每个算法都基于它必须处理的事件数量运行(Kafka 保留 30m)。目前,每个执行周期大约需要~90s,它在批处理模式下运行一个循环,如下所述:

  1. 从 Kafka 输入主题获取事件
  2. 围绕70 Spark SQL 处理
  3. 将输出发送回 Kafka 输出主题

由于每个周期大约需要 90s,这意味着 kafka 事件可以从90s 到 180s 来处理。我必须将此处理时间减少到60s

恕我直言,我可以扩展 spark 硬件以在批处理模式下寻找更好的 SQL 性能,但由于我确定算法处理的主要部分只是创建执行计划,我想知道可以对执行计划做什么肯定会显着减少处理时间。

它目前在Spark 3.0.1 上以20 vCore 服务器和32GB RAM 的独立配置运行。这是支持此问题的代码示例。希望这能解释这种情况。

代码示例

1。获取Kafka主题

streamdata = spark.read.format("kafka").option("kafka.bootstrap.servers", kafka_servers).\ 
                                                option("subscribe", readtopic).option("failOnDataLoss", "false").\ 
                                                option("startingOffsets", "earliest").\ 
                                                option("endingOffsets", "latest").load() 

 
#Besides columns identified bellow, we will bring from kafka : offset (auto number) and timestamp (insert datetime) 
streamdata = streamdata.withColumn('value', streamdata['value'].cast('string')).drop('key','topic','partition','timestampType') 

streamdata = streamdata.withColumn("LAST_UPDATE", split(col("value"), ",").getItem(0).cast(IntegerType()))\ 
                .withColumn("DESCRIPTION", split(col("value"), ",").getItem(1).cast(StringType())) 
.withColumn("MYTAG", split(col("value"), ",").getItem(6).cast(StringType())) 

2。窗口函数处理(SEVE$RAL QUERIES LIKE THAT OR MORE COMPLEX)

SUM_COUNT = spark.sql(""" 
    SELECT FIELD_A, FIELD_B, sum(FIELD_C) OVER (PARTITION BY FIELD_D ORDER BY CAST(LAST_UPDATE as timestamp) RANGE BETWEEN INTERVAL 12 HOURS PRECEDING AND CURRENT ROW) as FIELD_COUNT_12h  
    FROM streamdata 
""") 

3。将数据发送回 Kafka

#create value as a concat json of all columns 
SendKafka = query_03.withColumn("value", to_json(struct([query_03[x] for x in query_03.columns]))) 

#Send back to kafka 
SendKafka.write.format("kafka").option("kafka.bootstrap.servers", kafka_servers).option("topic", writetopic).save() 

【问题讨论】:

  • 如果你愿意分享更多关于转换的数据,也许我们可以试试 sqls。另外,集群的利用率是多少? (神经节报告可以派上用场)

标签: apache-spark pyspark apache-spark-sql databricks sql-execution-plan


【解决方案1】:

如果不查看实际执行计划,仅使用代码就很难说。从调优开始spark.sql.shuffle.partitions - 将其设置为可用于 Spark 作业的核心数 - 如果您有 20 个核心,并且使用等于 200 的默认值,则意味着在第一次 shuffle 之后,每个核心都会执行代码10 次 (200/20),前一次后一次(在 Spark 3 中,由于自适应查询执行,问题应该较少)。此外,考虑到 Spark 根据主题中的分区数从 Kafka 读取数据,因此如果您的分区数少于核心数,那么您的核心在读取时将处于空闲状态 - 检查 Kafka 连接器的minPartitions 选项(请参阅@987654321 @)

另外,请查看 Spark SQL tuning guideSpark Tuning guide 中的建议。

【讨论】:

  • 感谢亚历克斯的回答。让我试试你的建议,让你知道结果。我还在实施 Spark Streaming 以尝试作为我们解决方案的附加路线图。您对此也有任何线索吗?对于这种情况,您如何看待 Spark Streaming?非常感谢。
  • 对于结构化流,推荐与批处理重叠,但还有其他需要考虑的因素,例如检查点中的数据大小等。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-01-18
  • 1970-01-01
  • 2022-10-06
相关资源
最近更新 更多