You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

升级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一致性检查、系统表查询等,导致延迟。

疑问与求助

  1. 这些内部查询是否是性能下降的根源?版本间查询耗时增加的根本原因是什么?
  2. Insert、Update查询也出现耗时增加的情况,如何统一解决?
  3. 在不降级依赖的前提下,如何优化该性能问题?

内容的提问来源于stack exchange,提问作者Niranjana Datta

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 22:16:16