【问题标题】:how to use lag/lead function in spark streaming application?如何在火花流应用程序中使用滞后/领先功能?
【发布时间】: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


【解决方案1】:

我认为它不需要任何 leg 和 lead 函数。你也可以通过join 得到你想要的输出。检查以下代码以供参考:

注意:我在kafkaDF中添加了更多数据以便更多理解。

scala> kafkaDf.show(false)
+---------+----+-------+-----+----------+
|companyId|year|quarter|sales|prev_sales|
+---------+----+-------+-----+----------+
|15       |2016|4      |100.5|          |
|15       |2016|1      |115.8|          |
|15       |2016|3      |150.1|          |
+---------+----+-------+-----+----------+


scala> cassandraTabledf.show
+---------+----+-------+-----+----------+
|companyId|year|quarter|sales|prev_sales|
+---------+----+-------+-----+----------+
|       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|
+---------+----+-------+-----+----------+


scala>kafkaDf.alias("k").join(
                              cassandraTabledf.alias("c"), 
                              col("k.companyId") === col("c.companyId") && 
                              col("k.quarter") === col("c.quarter") && 
                              (col("k.year") - 1) === col("c.year"),
                              "left"
                             )
                       .drop("prev_sales")
                       .select(col("k.*"), col("c.sales").alias("prev_sales"))
                       .withColumn("prev_sales", when(col("prev_sales").isNull, col("sales")).otherwise(col("prev_sales")))
                       .show()
+---------+----+-------+-----+----------+
|companyId|year|quarter|sales|prev_sales|
+---------+----+-------+-----+----------+
|       15|2016|      1|115.8|     700.7|
|       15|2016|      3|150.1|     320.6|
|       15|2016|      4|100.5|     221.4|
+---------+----+-------+-----+----------+

【讨论】:

  • 这里的主要问题是这是一个 Spark Join,而不是仅针对必要数据的优化连接
  • @BdEngineer,我已经更改了连接条件,它将处理“如果 2015 年第 4 季度没有记录,它应该填充 2016 年第 4 季度的“销售额””的场景。请检查确认
  • 由于这个很接近,您可以在stackoverflow上提出新问题,并提供结果数据框和另一个具有价值的数据框。
  • @BdEngineer 我已回复,请查看stackoverflow.com/questions/59579922/…
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-06-23
  • 2018-08-09
  • 1970-01-01
  • 1970-01-01
  • 2015-11-07
  • 1970-01-01
相关资源
最近更新 更多