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

Spark Streaming与Cassandra Direct Join失效问题求助

问题排查:Spark Cassandra Connector DirectJoin失效与NoSuchMethodError问题

环境与核心需求

  • 技术栈:Spark 3.2.1、Cassandra 4.0.4、com.datastax.spark:spark-cassandra-connector_2.12:3.1.0
  • 数据源:Kafka
  • 核心需求:将Kafka消息转为DataFrame后,与Cassandra表(复合主键user_id+user_type)左连接,仅写入Cassandra中不存在的记录
  • 预期行为:利用SCC 2.5+支持的DataFrame DirectJoin(等价于RDD的joinWithCassandraTable),避免Spark端SortMergeJoin带来的性能损耗

问题现象

  1. DataSource V2场景:执行计划显示触发Spark端SortMergeJoin,未走Cassandra端的DirectJoin
  2. DataSource V1场景:设置spark.sql.cassandra.directJoinSetting=on强制开启DirectJoin后,抛出NoSuchMethodError报错

附相关信息

执行计划(DataSource V2)

== Physical Plan ==
SortMergeJoin [user_id#0, user_type#1], [user_id#10, user_type#11], LeftOuter
:- Sort [user_id#0 ASC, user_type#1 ASC], false, 0
:  +- Exchange hashpartitioning(user_id#0, user_type#1, 200), ENSURE_REQUIREMENTS, [id=#123]
:     +- Project ...(Kafka消息解析逻辑)
:        +- KafkaV2Relation ...
+- Sort [user_id#10 ASC, user_type#11 ASC], false, 0
   +- Exchange hashpartitioning(user_id#10, user_type#11, 200), ENSURE_REQUIREMENTS, [id=#456]
      +- CassandraTableScan ...

报错信息(DataSource V1)

java.lang.NoSuchMethodError: org.apache.spark.sql.cassandra.CassandraSourceRelation$.apply(Lorg/apache/spark/sql/SparkSession;Lorg/apache/spark/sql/catalyst/catalog/ExternalCatalogTable;Lscala/Option;Lscala/collection/immutable/Map;)Lorg/apache/spark/sql/cassandra/CassandraSourceRelation;
	at org.apache.spark.sql.cassandra.DefaultSource.createRelation(DefaultSource.scala:108)
	...

Spark提交命令

spark-submit \
  --class com.example.MySparkApp \
  --master yarn \
  --deploy-mode cluster \
  --packages com.datastax.spark:spark-cassandra-connector_2.12:3.1.0 \
  --conf spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions \
  --conf spark.sql.cassandra.directJoinSetting=on \
  --conf spark.sql.cassandra.source.useV1SourceList=* \
  my-app.jar

核心代码片段

// Kafka消息解析为DataFrame
val kafkaDF = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "xxx")
  .option("subscribe", "user-topic")
  .load()
  .selectExpr("CAST(value AS STRING)")
  .select(from_json(col("value"), userSchema).as("user"))
  .select("user.*")

// 读取Cassandra表
val cassandraDF = spark.read
  .format("org.apache.spark.sql.cassandra")
  .options(Map(
    "keyspace" -> "my_keyspace",
    "table" -> "user_table"
  ))
  .load()

// 左连接后过滤出不存在的记录
val newRecordsDF = kafkaDF.join(cassandraDF, 
  Seq("user_id", "user_type"), 
  "left_outer"
).filter(cassandraDF("user_id").isNull)

// 写入Cassandra
newRecordsDF.writeStream
  .format("org.apache.spark.sql.cassandra")
  .options(Map(
    "keyspace" -> "my_keyspace",
    "table" -> "user_table"
  ))
  .option("checkpointLocation", "/path/to/checkpoint")
  .start()
  .awaitTermination()

问题根源分析

1. DataSource V2未触发DirectJoin的原因

SCC 3.x的DataSource V2实现中,DirectJoin仅支持等值内连接(Inner Join),左外连接(Left Outer Join)不会触发DirectJoin,会退化为Spark端SortMergeJoin。原因是DirectJoin的底层逻辑是将主键过滤条件下推到Cassandra执行查询,而左外连接需要保留左表所有数据,无法完全下推到Cassandra完成。

2. DataSource V1出现NoSuchMethodError的原因

Spark 3.2.1与SCC 3.1.0的DataSource V1 API存在兼容性冲突:SCC 3.x的V1实现基于Spark 3.0+的API开发,但Spark 3.2对部分核心API的方法签名做了变更,导致反射调用时出现方法找不到的错误。spark.sql.cassandra.source.useV1SourceList=*强制所有Cassandra表使用V1数据源,进一步触发了这个兼容性问题。


解决方案

方案1:改用RDD API的joinWithCassandraTable(推荐)

既然DataFrame的DirectJoin不支持左外连接,直接用RDD API实现需求,该方式会将主键过滤下推到Cassandra,仅拉取存在的记录,性能最优:

import com.datastax.spark.connector._
import org.apache.spark.rdd.RDD

// KafkaDF转RDD,映射为(主键元组,原始记录)结构
val kafkaRDD = kafkaDF.rdd.map(row => 
  ((row.getAs[String]("user_id"), row.getAs[String]("user_type")), row)
)

// 与Cassandra表关联,过滤出不存在的记录
val newRecordsRDD = kafkaRDD.joinWithCassandraTable("my_keyspace", "user_table")
  .on(SomeColumns("user_id", "user_type"))
  .filter(_._2.isEmpty) // 仅保留Cassandra中无匹配的记录
  .map(_._1._2) // 取出原始Kafka记录

// 转DataFrame后写入Cassandra
spark.createDataFrame(newRecordsRDD, userSchema)
  .writeStream
  .format("org.apache.spark.sql.cassandra")
  .options(Map(
    "keyspace" -> "my_keyspace",
    "table" -> "user_table"
  ))
  .option("checkpointLocation", "/path/to/checkpoint")
  .start()
  .awaitTermination()

方案2:调整DataFrame逻辑减少Cassandra读取量

如果坚持使用DataFrame API,可先提取当前批次的主键列表,仅拉取Cassandra中匹配的记录,再做左连接:

// 提取当前批次的所有主键(去重)
val batchKeysDF = kafkaDF.select("user_id", "user_type").distinct()

// 仅拉取Cassandra中存在的对应主键记录
val existingCassandraDF = spark.read
  .format("org.apache.spark.sql.cassandra")
  .options(Map(
    "keyspace" -> "my_keyspace",
    "table" -> "user_table"
  ))
  .load()
  .join(batchKeysDF, Seq("user_id", "user_type"), "inner")

// 左连接过滤出新记录
val newRecordsDF = kafkaDF.join(existingCassandraDF, 
  Seq("user_id", "user_type"), 
  "left_outer"
).filter(existingCassandraDF("user_id").isNull)

注意:此方案需控制批次主键的数量,避免内存溢出。

方案3:降级Spark版本(不推荐)

若必须使用DataSource V1的DirectJoin,可将Spark版本降级到3.0.x,与SCC 3.1.0的V1实现完全兼容,但会丢失Spark 3.2的新特性。


内容的提问来源于stack exchange,提问作者Александр Трутнев

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 05:15:40