【发布时间】:2016-03-13 13:27:04
【问题描述】:
我有一种情况,我想使用 SparkSQL“迭代”或映射“宽行”而不是逻辑 Cassandra 行(CQL 行)。
基本上我的数据按timestamp(分区键)进行分区,并且有一个集群键,即传感器 ID。
对于每个timestamp我想做的操作,一个简单的例子就是做sensor1/sensor2。
我怎样才能通过保持数据局部性使用 SparkSQL 有效地做到这一点(而且我认为我的数据模型非常适合这些任务)?
我读到 this post on Datastax 在 Cassandra 连接器中提到了 spanBy 和 spanByKey。这将如何与 SparkSQL 一起使用?
伪代码示例(pySpark):
ds = sqlContext.sql("SELECT * FROM measurements WHERE timestamp > xxx")
# span the ds by clustering key
# filter the ds " sensor4 > yyy "
# for each wide-row do sensor4 / sensor1
【问题讨论】:
标签: apache-spark cassandra pyspark apache-spark-sql pyspark-sql