升级spark-cassandra-connector至3.1.0后性能骤降的排查与解决咨询
Spark-Cassandra Connector升级后查询性能骤降排查与优化求助
背景
我们维护着一个每日处理数百万消息的消息处理应用,基于Scala、Spark构建,使用Kafka与Cassandra DB,每条消息处理时需执行多笔DB查询。近期将应用从Spark集群迁移至K8s(使用mlrun、spark-operator),同步升级了Scala及依赖版本。
问题
迁移后发现Cassandra DB查询耗时显著增加,导致应用性能下降(消息处理变慢),而Cassandra集群未变更。
依赖版本对比
迁移前
<properties> <scala.tools.version>2.11</scala.tools.version> <scala.version>2.11.12</scala.version> <spark.version>2.4.0</spark.version> </properties> <dependency> <groupId>org.scala-lang</groupId> <artifactId>scala-library</artifactId> <version>${scala.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_${scala.tools.version}</artifactId> <version>${spark.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector-unshaded_2.11</artifactId> <version>2.4.0</version> </dependency>
迁移后
<properties> <scala.tools.version>2.12</scala.tools.version> <scala.version>2.12.10</scala.version> <spark.version>3.1.2</spark.version> </properties> <dependency> <groupId>org.scala-lang</groupId> <artifactId>scala-library</artifactId> <version>${scala.version}</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_${scala.tools.version}</artifactId> <version>${spark.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector_${scala.tools.version}</artifactId> <version>3.1.0</version> </dependency>
迁移变更点
- 升级Scala及相关依赖(如spark-cassandra-connector);
- 从Spark集群迁移至K8s,经本地测试证实该操作不影响查询耗时,性能差异源于依赖版本升级。
测试验证
为排查问题,编写了模拟业务逻辑的Scala测试程序:
import scala.io.Source import java.time.LocalDateTime import java.time.format.DateTimeFormatter import java.time.temporal.ChronoUnit import org.apache.spark.sql.SparkSession import com.datastax.spark.connector.toSparkContextFunctions import org.apache.log4j.{Level, LogManager} import org.joda.time.DateTime object SimpleCassandraProgram { def main(args: Array[String]): Unit = { val sBuild = SparkSession.builder() .appName(args(5)) .master(args(0).trim) // local .config("spark.cassandra.connection.host","<cassandra_host>") val spark = sBuild.getOrCreate() val sc = spark.sparkContext @transient lazy val log = LogManager.getLogger(sc.appName) LogManager.getRootLogger.setLevel(Level.DEBUG) val where_field = args(3) var count = 0 var totalTime: Long = 0 val fmt = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS") val txtFile = args(4).trim println(f"START TIME: " + DateTime.now()) // The file I used contains 10 lines => 10 select queries will be executed Source.fromFile(txtFile).getLines().foreach(line => { count += 1 var isProcessed = false val startTime = LocalDateTime.now log.info("Start of Query Execution: " + fmt.format(startTime) + " Run number: " + count) // args(1) is keyspace name, args(2) is table name val resRDD = sc.cassandraTable(args(1), args(2)).select("<field>").where(s"${where_field} =?", line) if (!resRDD.isEmpty()) { val latestRow = resRDD.first() val field_value = latestRow.getStringOption("<field>") if(field_value != None){ isProcessed=true log.info("Record has already been processed") } } val finishTime = LocalDateTime.now log.info("End of Query Execution: " + fmt.format(finishTime) + " Run number: " + count) val timeDiff = ChronoUnit.MILLIS.between(startTime, finishTime) log.info("Took " + timeDiff + "ms to Execute the Query.") // Excluding first query execution time since it includes the time to connect to Cassandra if (count != 1) totalTime += timeDiff }) println(f"FINISH TIME: " + DateTime.now()) println(f"AVERAGE QUERY EXECUTION TIME (EXCLUDING FIRST QUERY) - ${totalTime/(count - 1)}") } }
测试结果
- 迁移前依赖:平均查询耗时(排除首次连接)1468ms
- 迁移后依赖:平均查询耗时(排除首次连接)4109ms
通过调试日志发现,新版本依赖在执行业务查询前会执行大量内部查询,如Schema一致性检查、系统表查询等,导致延迟。
疑问与求助
- 这些内部查询是否是性能下降的根源?版本间查询耗时增加的根本原因是什么?
- Insert、Update查询也出现耗时增加的情况,如何统一解决?
- 在不降级依赖的前提下,如何优化该性能问题?
内容的提问来源于stack exchange,提问作者Niranjana Datta
相关产品推荐
相关产品推荐

