【问题标题】:Server side filtering of spark-cassandra on PySparkPySpark 上 spark-cassandra 的服务器端过滤
【发布时间】:2016-06-20 12:56:15
【问题描述】:

我是 Spark 的新手,想在与 Cassandra 合作时了解更多关于它的操作。

大部分教程都提醒我进行服务器端过滤,我完全理解这样做的重要性。

然而,这些教程要么基于 Scala,要么基于 pyspark_cassandra,而且都没有使用 PySpark。

只是好奇下面的 scriptlet 是否在进行服务器端过滤。

给定一个 SparkConf 对象conf

sc = pyspark.SparkContext(conf=conf)

sqlContext = SQLContext(sc)
df = (sqlContext.read.format("org.apache.spark.sql.cassandra")
    .options(keyspace="ks", table="tbl").load())

df.filter("id = 1234").show()

此外,在这种情况下,我是否将整个表加载到我的 spark 集群中进行过滤?

【问题讨论】:

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


    【解决方案1】:

    Cassandra 连接器支持 Spark DataFrames 上的谓词下推,因此只要启用下推,您就可以放心地假设基本过滤器在 Cassandra 端执行。它可能不适用于复杂的谓词。如果您有疑问,最好查看BasicCassandraPredicatePushDown docstrings

    您还可以检查执行计划(explain)。如果 predict 被下推,它应该列在PushedFilters 部分,例如:

    df = (sqlContext
      .read
      .format("org.apache.spark.sql.cassandra")
      .options(table="words", keyspace="test")
      .load())
    
    df.select("word").where(col("word") == "bar").explain()
    ## == Physical Plan ==
    ## Scan org.apache.spark.sql.cassandra.CassandraSourceRelation@62738171[word#0] 
    ## ... PushedFilters: [EqualTo(word,bar)]
    

    在 Spark 1.6 中,PushedFilters 的解释有点误导。它将列出数据源已显示的所有过滤器,但实际上不会告诉您数据源使用了哪些过滤器。在这种情况下,最好只查看explain 计划是否对谓词有单独的过滤步骤。如果确实如此,则连接器不会下推谓词。如果没有,则谓词被推送。

    另一个选项是为 Spark Cassandra 连接器打开 INFO/DEBUG 日志记录,以准确查看连接器在 Catalyst 中的作用

    【讨论】:

    • 我只想指出,在 Spark 1.6 上,PushedFilters 的解释具有误导性。它将列出数据源可以看到的所有过滤器,但实际上不会告诉您实际推送了哪些过滤器。在这种情况下,最好只查看 spark 是否在数据源之外执行了单独的“过滤”步骤。如果没有,则谓词被推送。您还可以打开连接器的 INFO/DEBUG 日志记录,以准确查看连接器在 Catalyst 中所做的事情。
    猜你喜欢
    • 2022-08-09
    • 2020-12-14
    • 1970-01-01
    • 2017-06-26
    • 2022-10-17
    • 2017-06-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多