【问题标题】:SELECT DISTINCT Cassandra in Spark在 Spark 中选择不同的 Cassandra
【发布时间】:2018-04-27 04:51:34
【问题描述】:

我需要一个查询,列出 spark 中唯一的复合分区键
CASSANDRA 中的查询:SELECT DISTINCT key1, key2, key3 FROM schema.table; 相当快,但是相比之下,将相同类型的数据过滤器放在 RDD 或 spark.sql 中检索结果的速度非常慢。

例如

---- SPARK ----
var t1 = sc.cassandraTable("schema","table").select("key1", "key2", "key3").distinct()
var t2 = spark.sql("SELECT DISTINCT key1, key2, key3 FROM schema.table")

t1.count // takes 20 minutes
t2.count // takes 20 minutes

---- CASSANDRA ----
// takes < 1 minute while also printing out all results
SELECT DISTINCT key1, key2, key3 FROM schema.table; 

表格格式如下:

CREATE TABLE schema.table (
    key1 text,
    key2 text,
    key3 text,
    ckey1 text,
    ckey2 text,
    v1 int,
    PRIMARY KEY ((key1, key2, key3), ckey1, ckey2)
);

不会在其查询中触发使用 cassandra 优化吗?
如何有效地检索这些信息?

【问题讨论】:

    标签: apache-spark cassandra distinct


    【解决方案1】:

    快速解答

    不会在其查询中触发使用 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 示例实际上将分为几个步骤。

    1. 扫描:读取该表中的所有数据。这意味着将每个值从 C 机器序列化到 Spark Executor JVM,换句话说,需要大量工作。
    2. *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,我们需要添加一个函数,例如已经存在的 whereselectlimit,调用 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 连接器,如上所述。

    【讨论】:

    • 一个示例或资源链接会很好。我真的不知道你在说什么。
    • 添加了关于 spark 的入门,它如何与 Cassandra 交互,以及问题中使用的方法如何发挥作用。
    • 感谢您在回答中添加了如此多的细节!我不认为我会触及“催化剂”:但很高兴知道 CassandraConnector 的存在。谷歌搜索继续......
    【解决方案2】:

    只要选择分区键,就可以使用CassandraRDD的.perPartitionLimit函数:

    val partition_keys = sc.cassandraTable("schema","table").select("key1", "key2", "key3").perPartitionLimit(1)
    

    这是有效的,因为根据SPARKC-436

    select key from some_table per partition limit 1

    给出与

    相同的结果

    select distinct key from some_table

    此功能是在 spark-cassandra-connector 2.0.0-RC1 中引入的 并且至少需要C* 3.6

    【讨论】:

    • 好答案!经过测试,快速准确。
    • 这也与 WHERE 子句兼容,这产生了一些有趣的用例;查找具有特定集群键值或值范围的所有分区键。
    【解决方案3】:

    Distinct 的性能很差。 这里有一些替代方案的好答案: How to efficiently select distinct rows on an RDD based on a subset of its columns`

    您可以使用 toDebugString 来了解您的代码混洗了多少数据。

    【讨论】:

      猜你喜欢
      • 2019-07-27
      • 2014-05-09
      • 2015-09-21
      • 2013-07-15
      • 2020-01-28
      • 2018-12-10
      • 2019-06-15
      • 2021-10-07
      • 2016-08-20
      相关资源
      最近更新 更多