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)大致相同。
与Spark Streaming不同,Spark本身正在尝试尽可能快地预加载尽可能多的数据,以便能够并行对其进行操作。因此,预加载是懒惰的,但触发时却很贪婪。但是,有cassandra扇区特定的因素:
调用限制将允许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 并比较结果。