【问题标题】:spark dataframe custom udf to return an arrayspark dataframe 自定义 udf 返回一个数组
【发布时间】:2018-07-20 06:36:01
【问题描述】:

单列样本数据集:

5.1,
4.9,
4.7,
4.6,
5,3.
5.4,
4.6,
5,
4.4,
4.9,
5.4,
4.8,
4.8,
4.3,
5.8

我希望它先按升序排序,然后间隔选择值并将其作为数组返回。

例如,如果区间 = 5,并且排序后的数据集是

4.3,
4.4,
4.6,
4.6,
4.7,
4.8,
4.8,
4.9,
4.9,
5,
5.1,
5,3.
5.4,
5.4,
5.8

它应该返回Array(4.3, 4.7, 5, 5.8)

有没有办法以乐观的方式做到这一点?

提前致谢 沙克蒂

这是我尝试过的,但无法获得第一个值。

val interval = 5
val count = df.count() //15
val n = (count/interval).toInt //3
println(s"interval: $interval, count: $count, n: $n")

val window = Window.orderBy("col1")
val sorted =  df.withColumn("rowId", functions.row_number().over(window))
sorted.show()

val sb = new StringBuilder
for (i <- 0 to n) {
  val intervalPoint = interval * i
  println(s"i: $i, intervalPoint: $intervalPoint")
  sb.append(s"rowId == $intervalPoint or ")
}

sb.delete(sb.size - 3, sb.size - 1)
println(s"sb: ${sb.toString()}") //rowId == 0 or rowId == 5 or rowId == 10 or rowId == 15

val intervals = sorted.where(sb.toString()).select("col1").collectAsList()
println(s"intervals: $intervals") //[[4.7], [5.0], [5.8]]

如您所见,首先它必须按 col 排序并附加一个行 ID。希望这两个可以在一次扫描中完成。并再次扫描整个数据集以获取间隔,我也无法获得第一个值。如果这必须应用于多个列,它必须在一个循环中,并且没有。扫描次数会增加。

【问题讨论】:

  • 向我们展示你到目前为止所做的尝试......
  • 您不需要 spark 并且 spark 无法满足您的要求。使用 scala 进行简单的本地计算应该是更好的选择
  • 这只是一个示例数据,实际数据可能有更多记录和更多列,我需要分别获取每个选定列的间隔。我不确定简单的 scala 是否能够处理大型数据集。我想使用数据框,如果需要 udfs

标签: sorting apache-spark dataframe intervals


【解决方案1】:

一种可能的解决方案(更清洁):

假设您有一个 数据框 'df' :

df.createOrReplaceTempView("temp_table_1")
sparkSession.sql("select col1 from (select ROW_NUMBER() OVER (ORDER BY col1) AS id, col1 from temp_table_1) y where id%(SOME_INTERVAL) = 0 order by col1").show()

【讨论】:

  • 当然,它更干净,但是对于多个列和大型数据集应用相同的方法呢?有没有办法减少扫描次数?
  • 我不确定您所说的不止一个列是什么意思?你能提供另一个数据集来满足你的要求吗?在上述方法中,子查询中对整个表的扫描次数为 1。既然这一切都在记忆中,我认为这并不重要。
猜你喜欢
  • 2016-12-24
  • 2017-11-01
  • 1970-01-01
  • 2017-05-11
  • 2019-02-26
  • 2021-11-26
  • 1970-01-01
相关资源
最近更新 更多