Spark查询唯一复合分区键性能远逊Cassandra的问题咨询
嘿,这个问题我太有发言权了!相信你也发现了,Cassandra对分区键的DISTINCT查询简直是天生优势,但Spark默认处理起来却慢得让人抓狂——其实核心原因就是两者的处理逻辑完全不一样,咱们针对性调整一下就能解决!
方案1:把查询推给Cassandra原生执行
Cassandra本身对分区键的DISTINCT查询是直接基于分区索引做的,根本不需要扫全表,所以速度飞快。那咱们直接让Spark把这个查询“甩”给Cassandra去执行就行了,不用自己费劲做全表扫描和去重。
用Spark SQL的方式
只需要先设置一个配置,开启Cassandra的DISTINCT查询下推:
spark.conf.set("spark.sql.cassandra.distinct.allowPushdown", "true")
然后再执行你的查询:
var t2 = spark.sql("SELECT DISTINCT key1, key2, key3 FROM schema.table") t2.count // 这下速度应该和Cassandra原生查询差不多了
用RDD的方式
如果习惯用RDD API,直接调用Cassandra的原生CQL查询:
val distinctKeys = sc.cassandraQuery("SELECT DISTINCT key1, key2, key3 FROM schema.table") .as((k1: String, k2: Int, k3: Long) => (k1, k2, k3)) distinctKeys.count
这个方式也是让Cassandra直接处理查询,结果返回给Spark,效率拉满。
方案2:直接读取Cassandra的系统元数据表
如果上面的配置不生效,或者你需要更底层的控制,还可以直接查Cassandra的系统表——它的元数据里已经存了所有分区的信息,根本不用碰业务数据。
Cassandra的system_schema.partitions表(旧版本是system.partitions)记录了所有表的分区详情,咱们可以直接查询这个表:
val partitionKeysDF = spark.sql(""" SELECT DISTINCT json_extract_scalar(partition_key, '$[0]') as key1, json_extract_scalar(partition_key, '$[1]') as key2, json_extract_scalar(partition_key, '$[2]') as key3 FROM system_schema.partitions WHERE keyspace_name = 'schema' AND table_name = 'table' """) partitionKeysDF.count
这里用json_extract_scalar是因为partition_key字段是JSON数组格式(对应复合分区键的各个字段),你可以根据自己的分区键类型调整解析方式。这个方法完全不需要扫描业务表,速度快到离谱。
为啥Spark默认这么慢?
最后给你唠唠原因:Spark的distinct操作是先把全表数据拉到集群各个节点,然后做shuffle、去重、合并——相当于把所有数据都过一遍,当然慢。而Cassandra的分区键是分布式存储的核心,每个节点只管理自己的分区,它可以直接从本地的分区索引里捞取唯一的分区键,然后汇总结果,全程不用扫全表,这差距能不大嘛!
内容的提问来源于stack exchange,提问作者ChiMo

