【问题标题】:How to aggregate data into ranges (bucketize)?如何将数据聚合到范围(桶化)?
【发布时间】:2017-04-12 06:02:04
【问题描述】:

我有一张这样的桌子

+---------------+------+
|id             | value|
+---------------+------+
|               1|118.0|
|               2|109.0|
|               3|113.0|
|               4| 82.0|
|               5| 60.0|
|               6|111.0|
|               7|107.0|
|               8| 84.0|
|               9| 91.0|
|              10|118.0|
+---------------+------+

ans 想将值聚合或合并到一个范围 0,10,20,30,40,...80,90,100,110,120我如何在 SQL 或更具体的 spark sql 中执行此操作?

目前我有一个横向视图加入范围,但这似乎相当笨拙/效率低下。

离散的分位数并不是我真正想要的,而是具有此范围的CUT

编辑

https://github.com/collectivemedia/spark-ext/blob/master/sparkext-mllib/src/main/scala/org/apache/spark/ml/feature/Binning.scala 会执行动态分箱,但我宁愿需要这个指定范围。

【问题讨论】:

  • org.apache.spark.ml.feature.Bucketizer 采用明确提供的分割点数组。然后,您应该能够在输出列上进行分组。
  • 我认为在这种情况下,建议的解决方案更简单/可能更有效。

标签: sql apache-spark apache-spark-sql


【解决方案1】:

一般情况下,可以使用org.apache.spark.ml.feature.Bucketizer进行静态分箱:

val df = Seq(
  (1, 118.0), (2, 109.0), (3, 113.0), (4, 82.0), (5, 60.0),
  (6, 111.0), (7, 107.0), (8,  84.0), (9, 91.0), (10, 118.0)
).toDF("id", "value")

val splits = (0 to 12).map(_ * 10.0).toArray

import org.apache.spark.ml.feature.Bucketizer
val bucketizer = new Bucketizer()
  .setInputCol("value")
  .setOutputCol("bucket")
  .setSplits(splits)

val bucketed = bucketizer.transform(df)

val solution = bucketed.groupBy($"bucket").agg(count($"id") as "count")

结果:

scala> solution.show
+------+-----+
|bucket|count|
+------+-----+
|   8.0|    2|
|  11.0|    4|
|  10.0|    2|
|   6.0|    1|
|   9.0|    1|
+------+-----+

当值位于定义的 bin 之外时,分桶器会抛出错误。可以将分割点定义为Double.NegativeInfinityDouble.PositiveInfinity 以捕获异常值。

Bucketizer 旨在通过对正确的存储桶执行二进制搜索来有效地处理任意拆分。对于像您这样的常规垃圾箱,您可以简单地执行以下操作:

val binned = df.withColumn("bucket", (($"value" - bin_min) / bin_width) cast "int")

其中bin_minbin_width 分别是最小 bin 的左侧区间和 bin 宽度。

【讨论】:

  • 但是假设一个桶是空的,分组也不会返回任何结果。因此,如果我想查看所有存储桶(以及计数为 0 的空存储桶)的列表,这可以在没有连接的情况下执行吗?
  • 在分箱后执行与范围的连接应该非常有效。
【解决方案2】:

用这个试试“GROUP BY”

SELECT id, (value DIV 10)*10 FROM table_name ;

以下将使用适用于 Scala 的数据集 API:

df.select(('value divide 10).cast("int")*10)

【讨论】:

    猜你喜欢
    • 2016-12-02
    • 1970-01-01
    • 1970-01-01
    • 2018-11-28
    • 1970-01-01
    • 2017-01-31
    • 1970-01-01
    • 2021-07-17
    • 2021-06-16
    相关资源
    最近更新 更多