Spark-Scala使用spark-cassandra-connector多表关联无法全部走Direct Join
问题根因
连续使用Spark原生join语法关联Cassandra表时,第一次join输出的DataFrame中,关联键id会丢失Cassandra Connector识别Direct Join所需的元标记,第二次join时连接器无法判定该id可作为分区键下推到Cassandra做点查,因此退化为全表扫描后执行SortMergeJoin。你单独关联任意一张表时基于原始流的id字段,所以可以正常触发Direct Join。
解决方案
方案1:使用joinWithCassandraTable专用API(推荐)
Cassandra Connector原生提供的joinWithCassandraTable专门为Cassandra关联做了优化,原生支持多次Direct Join,不会出现元数据丢失问题,代码示例如下:
import com.datastax.spark.connector._ val joinedDataDf = partitionKeyStream // 关联第一张Cassandra表 .joinWithCassandraTable("key1", "table1") .on(Columns("id") as "id") // 流的id字段对应Cassandra表的id分区键 .select("id", "latitude") // 关联第二张Cassandra表 .joinWithCassandraTable("key1", "table2") .on(Columns("id") as "id") .select("id", "latitude", "longitude")
方案2:拆分单表关联后合并
如果习惯用原生join语法,可分别基于原始流做两次单表关联后再合并结果,两次单表关联都可以触发Direct Join:
// 提取原始流的关联键字段 val baseIdStream = partitionKeyStream.select("id") // 分别做单表左关联,均触发Direct Join val joinTable1 = baseIdStream.join(cassandraTable1, baseIdStream("id") === cassandraTable1("id"), "left") val joinTable2 = baseIdStream.join(cassandraTable2, baseIdStream("id") === cassandraTable2("id"), "left") // 内存中合并两个关联结果 val joinedDataDf = joinTable1.join(joinTable2, Seq("id"), "left")
必要配置校验
请确认Spark配置中已经开启了Cassandra扩展,这是3.x版本支持Direct Join的前提:
spark.sql.extensions = com.datastax.spark.connector.CassandraSparkExtensions
同时需要确认两次关联的id字段都是对应Cassandra表的分区键,只有分区键才能触发Direct Join优化。
内容的提问来源于stack exchange,提问作者ilmar
相关产品推荐
相关产品推荐

