【问题标题】:Spark dataframe data aggregationSpark 数据帧数据聚合
【发布时间】:2017-05-12 07:17:09
【问题描述】:

我有以下要求在 scala 中聚合 Spark 数据帧上的数据。

我有一个包含两列的 spark 数据框。

mo_id   sales
201601  11.01
201602  12.01
201603  13.01
201604  14.01
201605  15.01
201606  16.01
201607  17.01
201608  18.01
201609  19.01
201610  20.01
201611  21.01
201612  22.01

如上所示,数据框有两列“mo_id”和“sales”。 我想在数据框中添加一个新列 (agg_sales),该列应包含截至当前月份的销售额总和,如下所示。

mo_id   sales   agg_sales
201601  10  10
201602  20  30
201603  30  60
201604  40  100
201605  50  150
201606  60  210
201607  70  280
201608  80  360
201609  90  450
201610  100 550
201611  110 660
201612  120 780

说明:

对于 201603 月份,agg_sales 将是从 201601 到 201603 的销售额总和。 对于 201604 月份,agg_sales 将是从 201601 到 201604 的销售额总和。 等等。

任何人都可以帮忙吗?

使用的版本:Spark 1.6.2 和 Scala 2.10

【问题讨论】:

  • 您的意思是要将sales 格式化为第一个数据集还是第二个数据集?
  • 我有一个包含两列的第一个数据框。
  • 所以在下一个数据框中我想添加一个新列(agg_sales)。
  • 所以在新数据集中,我总共有 3 列。 (month_id, sales, agg_sales)
  • 在下面查看我的答案

标签: scala apache-spark


【解决方案1】:

您正在寻找可以通过窗口函数完成的累积和:

scala> val df = sc.parallelize(Seq((201601, 10), (201602, 20), (201603, 30), (201604, 40), (201605, 50), (201606, 60), (201607, 70), (201608, 80), (201609, 90), (201610, 100), (201611, 110), (201612, 120))).toDF("id","sales")
df: org.apache.spark.sql.DataFrame = [id: int, sales: int]

scala> import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.expressions.Window

scala> val ordering = Window.orderBy("id")
ordering: org.apache.spark.sql.expressions.WindowSpec = org.apache.spark.sql.expressions.WindowSpec@75d454a4

scala> df.withColumn("agg_sales", sum($"sales").over(ordering)).show 
16/12/27 21:11:35 WARN WindowExec: No Partition Defined for Window operation! Moving all data to a single partition, this can cause serious performance degradation.
+------+-----+-------------+
|    id|sales|  agg_sales  |
+------+-----+-------------+
|201601|   10|           10|
|201602|   20|           30|
|201603|   30|           60|
|201604|   40|          100|
|201605|   50|          150|
|201606|   60|          210|
|201607|   70|          280|
|201608|   80|          360|
|201609|   90|          450|
|201610|  100|          550|
|201611|  110|          660|
|201612|  120|          780|
+------+-----+-------------+

请注意,我在ids 上定义了ordering,您可能需要某种时间戳来排序总和。

【讨论】:

  • 谢谢。让我试试。
  • 没有为窗口操作定义分区!将所有数据移动到单个分区
  • 只是一个警告,无需担心(如果您的数据大到可以担心分区,您可能已经定义了一个)
  • 是的,在我正在处理的实际数据集中,有一个我可以选择作为分区键的列。谢谢
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-10-09
  • 2018-03-04
  • 1970-01-01
  • 2018-09-03
  • 1970-01-01
  • 2021-11-26
  • 2015-05-19
相关资源
最近更新 更多