【发布时间】:2020-04-20 19:56:40
【问题描述】:
我使用 spark-sql 2.4.x 版本,datastax-spark-cassandra-connector 用于 Cassandra-3.x 版本。和卡夫卡一起。
我有一个来自 kafka 主题的财务数据的场景。 比如companyId、year、 Quarter、sales、prev_sales数据。
val kafkaDf = sc.parallelize(Seq((15,2016, 4, 100.5,"")).toDF("companyId", "year","quarter", "sales","prev_sales")
我需要使用来自 cassandra 表的上一年同季度数据进行 prev_sales,如下所示
val cassandraTabledf = sc.parallelize(Seq(
(15,2016, 3, 120.6, 320.6),
(15,2016, 2, 450.2,650.2),
(15,2016, 1, 200.7,700.7),
(15,2015, 4, 221.4,400),
(15,2015, 3, 320.6,300),
(15,2015, 2, 650.2,200),
(15,2015, 1, 700.7,100))).toDF("companyId", "year","quarter", "sales","prev_sales")
即对于 Seq((15,2016, 4, 100.5,"") 数据,它应该是 2015 年第 4 季度数据,即 221.4
所以新数据是
(15,2016, 4, 100.5,221.4)
如何做/实现这一目标? 我们可以进行显式查询,但是有什么方法可以在 cassandra 表上使用 join 来使用“滞后”功能?
【问题讨论】:
-
我不熟悉 Cassandra Spark 驱动程序,抱歉。任何 spark sql 查询都可能出现延迟,AFAIK
-
@BdEngineer,您能否详细说明您是如何加入数据框的?为什么它需要腿部功能?
标签: apache-spark cassandra apache-spark-sql