【问题标题】:Spark SQL and Cassandra JOINSpark SQL 和 Cassandra JOIN
【发布时间】:2016-05-25 08:51:05
【问题描述】:

我的 Cassandra 架构包含一个表,其分区键是时间戳,parameter 列是集群键。

每个分区包含 10k+ 行。这是以每秒 1 个分区的速度记录数据。

另一方面,用户可以定义“数据集”,我有另一个表,其中包含作为分区键的“数据集名称”和一个集群列,该列是引用另一个表的时间戳(因此是“数据集”是一个分区键列表)。

当然,我想做的事情看起来像是 Cassandra 的反模式,因为我想加入两个表。

但是使用 Spark SQL,我可以运行这样的查询并执行 JOIN

SELECT * from datasets JOIN data 
    WHERE data.timestamp = datasets.timestamp AND datasets.name = 'my_dataset'

现在的问题是:Spark SQL 是否足够智能以仅读取与datasets 中定义的timestamps 对应的data 的分区?

【问题讨论】:

    标签: apache-spark cassandra apache-spark-sql


    【解决方案1】:

    编辑:修正关于连接优化的答案

    Spark SQL 是否足够智能以仅读取与数据集中定义的时间戳相对应的数据分区?

    没有。实际上,由于您为数据集表提供了分区键,Spark/Cassandra 连接器将执行 谓词下推 并直接在 Cassandra 中使用 CQL 执行分区限制。但是连接操作本身不会有谓词下推,除非您使用带有joinWithCassandraTable()的RDD API

    请参阅此处了解所有可能的谓词下推情况:https://github.com/datastax/spark-cassandra-connector/blob/master/spark-cassandra-connector/src/main/scala/org/apache/spark/sql/cassandra/BasicCassandraPredicatePushDown.scala

    【讨论】:

    • 你确定它可以推送column = column形式的谓词吗?如果是这样,您可以提供一些参考。从我目前所见,Spark 只考虑column = value 形式的谓词。
    • 连接没有优化,但 AND datasets.name = 'my_dataset' 有谓词下推。如果您希望 spark/cassandra 连接器优化连接,您需要使用 RDD 编程 API (joinWithCassandraTable)
    • 谢谢。所以答案应该是否定的,不是吗?据我了解,OP 询问的是连接条件而不是谓词。
    • 感谢您的回答和 cmets!所以这意味着我应该添加类似AND data.timestamp IN (x) 的内容,其中x 是通过读取我的datasets 表中的分区获得的时间戳列表。这有意义吗?
    猜你喜欢
    • 2019-04-10
    • 1970-01-01
    • 2016-09-15
    • 1970-01-01
    • 2015-11-10
    • 2020-05-25
    • 1970-01-01
    • 2017-08-27
    • 1970-01-01
    相关资源
    最近更新 更多