【问题标题】:How to execute CQL query using pyspark如何使用 pyspark 执行 CQL 查询
【发布时间】:2020-11-12 04:01:14
【问题描述】:
我想使用 PySpark 执行 Cassandra CQL 查询。但我没有找到执行它的方法。我可以将整个表加载到数据框并创建 Tempview 并对其进行查询。
df = spark.read.format("org.apache.spark.sql.cassandra").
options(table="country_production2",keyspace="country").load()
df.createOrReplaceTempView("Test")
请建议任何更好的方法,以便我可以在 PySpark 中执行 CQL 查询。
【问题讨论】:
标签:
apache-spark
pyspark
cassandra
spark-cassandra-connector
【解决方案1】:
Spark SQL 不直接支持 Cassandra 的 cql 方言。它只允许您将表加载为 Dataframe 并对其进行操作。
如果您担心读取整个表来查询它,那么您可以使用下面给出的过滤器让 Spark 推送谓词以仅加载您需要的数据。
from pyspark.sql.functions import *
df = spark.read\
.format("org.apache.spark.sql.cassandra")\
.options(table=table_name, keyspace=keys_space_name)\
.load()\
.filter(col("id")=="A")
df.createOrReplaceTempView("Test")
【解决方案2】:
在 pyspark 中,您使用的是 SQL,而不是 CQL。如果 SQL 查询以某种方式与 CQL 匹配,即您正在按分区或主键查询,那么 Spark Cassandra 连接器 (SCC) 会将查询转换为该 CQL,并执行(所谓的谓词下推)。如果不匹配,则 Spark 将通过 SCC 加载所有数据,并在 Spark 级别进行过滤。
所以在你注册了临时视图之后,你可以这样做:
val result = spark.sql("select ... from Test where ...")
并使用result 变量中的结果。要检查谓词下推是否发生,请执行result.explain(),并检查PushedFilters 部分的条件中的* 标记。