快速解答
不会在其查询中触发使用 cassandra 优化吗?
是的。但是使用 SparkSQL 只有列修剪和谓词下推。在 RDD 中是手动的。
我怎样才能有效地检索这些信息?
由于您的请求返回的速度足够快,我将直接使用 Java 驱动程序来获取此结果集。
长答案
虽然 Spark SQL 可以提供一些基于 C* 的优化,但这些优化通常仅限于使用 DataFrame 接口时的谓词下推。这是因为框架仅向数据源提供有限的信息。我们可以通过对您编写的查询进行说明来看到这一点。
让我们从 SparkSQL 示例开始
scala> spark.sql("SELECT DISTINCT key1, key2, key3 FROM test.tab").explain
== Physical Plan ==
*HashAggregate(keys=[key1#30, key2#31, key3#32], functions=[])
+- Exchange hashpartitioning(key1#30, key2#31, key3#32, 200)
+- *HashAggregate(keys=[key1#30, key2#31, key3#32], functions=[])
+- *Scan org.apache.spark.sql.cassandra.CassandraSourceRelation test.tab[key1#30,key2#31,key3#32] ReadSchema: struct<key1:string,key2:string,key3:string>
因此,您的 Spark 示例实际上将分为几个步骤。
-
扫描:读取该表中的所有数据。这意味着将每个值从 C 机器序列化到 Spark Executor JVM,换句话说,需要大量工作。
- *HashAggregate/Exchange/Hash Aggregate:从每个执行器获取值,在本地对其进行哈希处理,然后在机器之间交换数据并再次哈希以确保唯一性。通俗地说,这意味着创建大型哈希结构,将它们序列化,运行复杂的分布式排序合并,然后运行
再次哈希。 (昂贵)
为什么不将其中的任何一个推到 C* 中?这是因为Datasource(在这种情况下为 CassandraSourceRelation)没有提供有关查询的 Distinct 部分的信息。这只是 Spark 当前工作方式的一部分。 Docs on what is pushable
那么RDD版本呢?
通过 RDDS,我们向 Spark 提供了一组直接指令。这意味着如果你想向下推一些东西,它必须是manually specified。我们来看看RDD请求的调试输出
scala> sc.cassandraTable("test","tab").distinct.toDebugString
res2: String =
(13) MapPartitionsRDD[7] at distinct at <console>:45 []
| ShuffledRDD[6] at distinct at <console>:45 []
+-(13) MapPartitionsRDD[5] at distinct at <console>:45 []
| CassandraTableScanRDD[4] at RDD at CassandraRDD.scala:19 []
这里的问题是您的“不同”调用是对 RDD 的通用操作,而不是特定于 Cassandra。由于 RDD 要求所有优化都是显式的(你输入的就是你得到的),Cassandra 从来没有听说过这种对“不同”的需求,我们得到的计划几乎与我们的 Spark SQL 版本相同。进行全面扫描,将 Cassandra 中的所有数据序列化到 Spark。进行随机播放,然后返回结果。
那么我们能做些什么呢?
使用 SparkSQL,这与我们在不向 Catalyst(SparkSQL/Dataframes 优化器)添加新规则的情况下获得的性能差不多,让它知道 Cassandra 可以处理服务器上的一些不同调用等级。然后需要为 CassandraRDD 子类实现它。
对于 RDD,我们需要添加一个函数,例如已经存在的 where、select 和 limit,调用 Cassandra RDD。可以添加一个新的Distinct 调用here,尽管它只在特定情况下是允许的。这是一个目前在 SCC 中不存在的函数,但可以相对容易地添加,因为它所做的只是将 DISTINCT 添加到 requests 并可能添加一些检查以确保它是有意义的 DISTINCT。
我们现在可以在不修改底层连接器的情况下做什么?
由于我们知道我们想要发出的确切 CQL 请求,我们始终可以直接使用 Cassandra 驱动程序来获取此信息。 Spark Cassandra 连接器提供了一个我们可以使用的驱动程序池,或者我们可以直接使用 Java 驱动程序。要使用池,我们会执行类似的操作
import com.datastax.spark.connector.cql.CassandraConnector
CassandraConnector(sc.getConf).withSessionDo{ session =>
session.execute("SELECT DISTINCT key1, key2, key3 FROM test.tab;").all()
}
如果需要进一步的 Spark 工作,则将结果并行化。如果我们真的想分发它,很可能有必要将函数添加到 Spark Cassandra 连接器,如上所述。