【发布时间】:2021-07-04 08:21:08
【问题描述】:
我刚刚开发了一个 Spark SQL 应用程序,在一些算法分析过程中,我意识到执行计划需要大量时间来处理。如何优化 Spark SQL 执行计划的性能?
我查看了我们社区中有关此问题的几个问题/答案,但在我看来,没有任何内容能直截了当地执行此操作。因此,我希望得到一些社区支持来克服我的障碍,并可能在此过程中为进一步的开发人员留下路线图。
以下是有关所做工作的一些详细信息。
我开发了一个 Spark 应用程序,它会定期从 kafka 获取事件,对其进行处理并将输出再次发送回 kafka。简而言之,Spark 算法过滤/丰富信息并针对每个事件执行繁重而复杂的窗口滞后函数。
Spark 算法循环运行,因此每个算法都基于它必须处理的事件数量运行(Kafka 保留 30m)。目前,每个执行周期大约需要~90s,它在批处理模式下运行一个循环,如下所述:
- 从 Kafka 输入主题获取事件
- 围绕
70Spark SQL 处理 - 将输出发送回 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