使用数据帧时如何下推Cassandra的限制谓词?

add*_*ing 6 scala cassandra spark-cassandra-connector spark-dataframe

我有大卡桑德拉桌子。我只想从Cassandra加载50行。以下代码

val ds = sparkSession.read
      .format("org.apache.spark.sql.cassandra")
      .options(Map("table" -> s"$Aggregates", "keyspace" -> s"$KeySpace"))
      .load()
      .where(col("aggregate_type") === "DAY")
      .where(col("start_time") <= "2018-03-28")
      .limit(50).collect()
Run Code Online (Sandbox Code Playgroud)

以下代码将两个谓词从where方法中推入,但不限制一个。是否获取了整个数据(100万条记录)?如果不是,为什么此代码的运行时间与没有代码的代码limit(50)大致相同。

dk1*_*k14 5

与Spark Streaming不同,Spark本身正在尝试尽可能快地预加载尽可能多的数据,以便能够并行对其进行操作。因此,预加载是懒惰的,但触发时却很贪婪。但是,有cassandra扇区特定的因素:

  • 自动谓语下推有效 “其中”条款。

  • 根据此答案 limit(...),不会转换为CQL的答案LIMIT,因此其行为取决于下载足够数据后创建了多少个提取作业。引用:

调用限制将允许Spark跳过从基础DataSource读取某些部分的操作。这些将通过取消执行任务来限制从Cassandra读取的数据量。

可能的解决方案:

  • DataFrame限制可以通过限制numPartitions和数据交换速率(concurrent.reads和其他参数)来部分管理。如果您在大多数情况下都可以接受n〜50,那么您也可以限制where(dayIndex < 50 * factor * num_records)

  • 有一种LIMIT通过设置CQL的方法SparkPartitionLimit,它直接影响每个CQL请求(请参阅更多信息)-请记住,请求是按火花分割的。它在CassandraRdd扩展类中可用,因此您必须先转换为RDD。

该代码将类似于:

filteredDataFrame.rdd.asInstanceOf[CassandraRDD].limit(n).take(n).collect()
Run Code Online (Sandbox Code Playgroud)

这将附加LIMIT $N到每个CQL请求中。与DataFrame的限制不同,如果您limit多次指定CassandraRDD (.limit(10).limit(20))-只会附加最后一个。另外,我使用n而不是使用n / numPartitions + 1它(即使Spark和Cassandra分区是一对一的),每个分区返回的结果也会更少。结果,我必须添加take(n)才能<= numPartitions * n减小到n

警告where再次确认您的可以翻译成CQL(使用explain()),否则LIMIT将在过滤之前应用。

PS您还可以尝试使用sparkSession.sql(...)如此处)直接运行CQL 并比较结果。