当我们使用spark从csv for DB读取数据时,如下所示,它会自动将数据拆分到多个分区并发送到执行器
spark
.read
.option("delimiter", ",")
.option("header", "true")
.option("mergeSchema", "true")
.option("codec", properties.getProperty("sparkCodeC"))
.format(properties.getProperty("fileFormat"))
.load(inputFile)目前,我有一个id列表:
[1,2,3,4,5,6,7,8,9,...1000]我想要做的是将这个列表分割成多个分区,并发送到每个executor,在每个executor中运行sql
ids.foreach(id => {
select * from table where id = id
})当我们从cassandra加载数据时,连接器将生成如下查询sql:
select columns from table where Token(k) >= ? and Token(k) <= ? 这意味着,连接器将扫描整个数据库,实际上,我不需要扫描整个表,我只需要从id列表中的k(分区键)的表中获取所有数据。
表模式为:
CREATE TABLE IF NOT EXISTS tab.events (
k int,
o text,
event text
PRIMARY KEY (k,o)
);或者,我如何使用spark使用预定义的sql语句从cassandra加载数据,而无需扫描整个表?
发布于 2019-02-05 04:27:31
您只需使用joinWithCassandra function执行选择,只需选择您的操作所需的数据。但请注意,此功能只能通过RDD API使用。
如下所示:
val joinWithRDD = your_df.rdd.joinWithCassandraTable("tab","events")您需要确保DataFrame中的列名与Cassandra中的分区键名相匹配-有关更多信息,请参阅文档。
DataFrame实现仅在following blog post中描述的Spark Cassandra Connector的DSE版本中可用。
2020年9月更新:在Spark Cassandra Connector 2.5.0中添加了对加入Cassandra的支持
https://stackoverflow.com/questions/54523004
复制相似问题