【问题标题】:Spark window custom function - getting the total number of partition recordsSpark窗口自定义函数——获取分区记录总数
【发布时间】:2019-01-31 17:18:05
【问题描述】:

我有一个时间序列数据集,它按 id 分区,并按时间戳排序。示例:

  ID     Timestamp   Feature

 "XSC"   1986-05-21  44.7530
 "XSC"   1986-05-22  44.7530
 "XSC"   1986-05-23  23.5678

 "TM"    1982-03-08  22.2734
 "TM"    1982-03-09  22.1941
 "TM"    1982-03-10  22.0847
 "TM"    1982-03-11  22.1741
 "TM"    1982-03-12  22.1840
 "TM"    1982-03-15  22.1344

我有一些我需要计算的自定义逻辑,它应该在每个窗口、每个分区内完成。 我知道 Spark 对窗口函数有丰富的支持,我正在尝试将其用于此目的。

我的逻辑需要当前窗口/分区中的元素总数,作为标量。我需要它来做一些特定的计算(基本上,一个 for 循环最多)。

我尝试添加一个计数列,方法是

val window = Window.partitionBy("id").orderBy("timestamp") 
frame = frame.withColumn("my_cnt", count(column).over(window))

我需要做类似的事情:

var i = 1
var y = col("Feature")
var result = y
while (i < /* total number of records within each partition goes here */) {
    result = result + lit(1) * lag(y, i).over(window) + /* complex computation */
    i = i + 1
}
dataFrame.withColumn("Computed_Value", result)

如何将每个分区中的记录总数作为标量值?我还添加了计数“my_cnt”值,它添加了分区的总值,但在我的情况下似乎无法使用它。

【问题讨论】:

  • 你能展示一些示例输入和预期的输出吗?不清楚您要做什么?
  • 添加了输入和一些示例代码。
  • 还是不清楚。根据您的输入清楚地提供预期的输出。
  • 也许你需要一个聚合窗口函数,比如blog.nuvola-tech.com/2017/10/…
  • @1pluszara - 这里的问题不在于输出应该是什么。没关系。重要的是,我如何访问当前窗口/分区中的元素总数。粘贴在那里的代码只是一些逻辑来查看我需要总计数的位置,以及我需要它的格式(作为实际值,而不是列)

标签: apache-spark apache-spark-sql


【解决方案1】:

Spark 的collect_list 函数允许您将窗口值聚合为一个列表。可以将此列表传递给udf 以进行一些复杂的计算

如果你有来源

val data = List(
  ("XSC", "1986-05-21", 44.7530),
  ("XSC", "1986-05-22", 44.7530),
  ("XSC", "1986-05-23", 23.5678),
  ("TM", "1982-03-08", 22.2734),
  ("TM", "1982-03-09", 22.1941),
  ("TM", "1982-03-10", 22.0847),
  ("TM", "1982-03-11", 22.1741),
  ("TM", "1982-03-12", 22.1840),
  ("TM", "1982-03-15", 22.1344),
).toDF("id", "timestamp", "feature")
  .withColumn("timestamp", to_date('timestamp))

还有一些复杂的函数,包裹在你记录中的 UDF 中(例如表示为元组)

 val complexComputationUDF = udf((list: Seq[Row]) => {
  list
    .map(row => (row.getString(0), row.getDate(1).getTime, row.getDouble(2)))
    .sortBy(-_._2)
    .foldLeft(0.0) {
      case (acc, (id, timestamp, feature)) => acc + feature
    }
})

您可以定义将所有分区数据传递给每条记录的窗口,或者在有序窗口的情况下,将增量数据传递给每条记录

val windowAll = Window.partitionBy("id")
val windowRunning = Window.partitionBy("id").orderBy("timestamp")

并将它们放在一个新的数据集中,例如:

val newData = data
  // I assuming thatyou need id,timestamp & feature for the complex computattion. So I create a struct
  .withColumn("record", struct('id, 'timestamp, 'feature))
  // Collect all records in the partition as a list of tuples and pass them to the complexComupation
  .withColumn("computedValueAll",
     complexComupationUDF(collect_list('record).over(windowAll)))
  // Collect records in a time ordered windows in the partition as a list of tuples and pass them to the complexComupation
  .withColumn("computedValueRunning",
     complexComupationUDF(collect_list('record).over(windowRunning)))

这将导致类似:

+---+----------+-------+--------------------------+------------------+--------------------+
|id |timestamp |feature|record                    |computedValueAll  |computedValueRunning|
+---+----------+-------+--------------------------+------------------+--------------------+
|XSC|1986-05-21|44.753 |[XSC, 1986-05-21, 44.753] |113.07379999999999|44.753              |
|XSC|1986-05-22|44.753 |[XSC, 1986-05-22, 44.753] |113.07379999999999|89.506              |
|XSC|1986-05-23|23.5678|[XSC, 1986-05-23, 23.5678]|113.07379999999999|113.07379999999999  |
|TM |1982-03-08|22.2734|[TM, 1982-03-08, 22.2734] |133.0447          |22.2734             |
|TM |1982-03-09|22.1941|[TM, 1982-03-09, 22.1941] |133.0447          |44.4675             |
|TM |1982-03-10|22.0847|[TM, 1982-03-10, 22.0847] |133.0447          |66.5522             |
|TM |1982-03-11|22.1741|[TM, 1982-03-11, 22.1741] |133.0447          |88.7263             |
|TM |1982-03-12|22.184 |[TM, 1982-03-12, 22.184]  |133.0447          |110.91029999999999  |
|TM |1982-03-15|22.1344|[TM, 1982-03-15, 22.1344] |133.0447          |133.0447            |
+---+----------+-------+--------------------------+------------------+--------------------+

【讨论】:

  • 信息量大、描述性强且总体上很好的答案。它让我走上了正确的轨道。谢谢!
猜你喜欢
  • 2018-02-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-03-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多