Spark Streaming与Cassandra Direct Join失效问题求助
环境与核心需求
- 技术栈: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带来的性能损耗
问题现象
- DataSource V2场景:执行计划显示触发Spark端
SortMergeJoin,未走Cassandra端的DirectJoin - 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,提问作者Александр Трутнев

