【问题标题】:Scala RDD - Relaxing data aggregation based on criteriaScala RDD - 基于标准放宽数据聚合
【发布时间】:2018-01-10 03:29:41
【问题描述】:

我有这种格式的数据:

RDD[(String, String, String), Int)]

我可以这样表示

|------|------|------------|----------------|
|(Month|Gender|Nationality)|NumberOfCustomer|
|------|------|------------|----------------|
| 1    |  M   |    FRA     |      8         |
| 1    |  F   |    FRA     |      2         |
| 1    |      |    FRA     |      2         |
| 1    |  M   |            |      7         |
| 1    |  F   |            |      2         |
| 1    |  M   |    USA     |      3         |
| 1    |  F   |    USA     |      4         |
| 1    |      |    USA     |      13        |
|------|------|------------|----------------|

由于某些限制,当客户少于 10 个时,我无法显示数据。 因此,我需要通过放宽一些标准(nationality 然后 gender)来汇总数据。

例如,由于Month 1 and the Gender M and the Nationality FRA 没有足够的(少于 10 个)客户,我需要将数据连接到其他国籍(未知)。 处理数据后,我应该有这样的东西:

|------|------|------------|----------------|
|(Month|Gender|Nationality)|NumberOfCustomer|
|------|------|------------|----------------|
| 1    |  M   |  Other     |      15        |
|------|------|------------|----------------|

Month 1, Gender F and Nationality FRA 也是如此,然后是美国等等。 结果应该是:

|------|------|------------|----------------|
|(Month|Gender|Nationality)|NumberOfCustomer|
|------|------|------------|----------------|
| 1    |      |    FRA     |      2         |
| 1    |  M   |    Other   |      18        |
| 1    |  F   |    Other   |      8         |
| 1    |      |    USA     |      13        |
|------|------|------------|----------------|

之后,Month 1, Gender Unknown and Nationality FRA 的客户仍然不足。 所以我需要将它与Month 1, Gender Unknown and the Nationality 与(这里美国)最少的客户连接起来 结果:

|------|------|------------|----------------|
|(Month|Gender|Nationality)|NumberOfCustomer|
|------|------|------------|----------------|
| 1    |  M   |    Other   |      18        |
| 1    |  F   |    Other   |      8         |
| 1    |      |    USA     |      15        |
|------|------|------------|----------------|

之后,Month 1, Gender F and Nationality Other 的客户仍然不够。

我需要保留Month 1 - Gender Unknown - USA Nationality(因为有超过10个客户) 但我需要将国籍标准删除为Month 1 - Gender F - Other Nationality,因为其中只有 8 个客户。

最终的最终结果应该是

|------|------|------------|----------------|
|(Month|Gender|Nationality)|NumberOfCustomer|
|------|------|------------|----------------|
| 1    |Other |    Other   |      26        |
| 1    |      |    USA     |      15        |
|------|------|------------|----------------|

我的问题是如何使用 Apache Spark 中的 Scala RDD 尽可能高效地实现这一点? (放宽2个以上的标准,比如国籍,然后是性别,然后是年龄,然后是体重等等,总是以相同的顺序)

编辑:按照评论中的要求添加代码和数据

获取我的 RDD[(String, String, String), Int)] :

val reducedByKey = myRDD.map(x =>
  (
    (
      x.month,
      x.gender,
      x.nationality
    ), 1
  )
).reduceByKey(_+_)

一些数据:

((1,M,FRA),8)
((1,F,FRA),4)
((1,,FRA),46)
((1,M,ENG),13)
((1,F,ENG),40)
((1,M,USA),1)
((1,F,USA),4)
((1,,USA),3)
((2,M,FRA),4)
((2,F,FRA),1)
((2,M,USA),10)
((2,F,USA),4)
((2,,USA),60)

【问题讨论】:

  • 您能否添加一些数据和代码,以便您的问题可以重现?
  • @tuxdna 我编辑了问题
  • 为什么第一次聚合后客户数是15?
  • @RameshMaharjan 因为我将所有具有 FRA 国籍的男性与所有具有未知国籍 (8+7) 的男性合并
  • 我已经猜到了,但真正的问题是乳清 7 没有与美国的 3 合并? 7 与 FRA 的 8 合并的逻辑是什么?

标签: scala apache-spark rdd


【解决方案1】:

您的问题可以概括为NumberOfCustomer 列小于10 时,将GenderNationality 列更改为Other

因此,如果您知道将rdd 转换为dataframe,如下所示

+-----+------+-----------+----------------+
|Month|Gender|Nationality|NumberOfCustomer|
+-----+------+-----------+----------------+
|    1|     M|        FRA|               8|
|    1|     F|        FRA|               4|
|    1|      |        FRA|              46|
|    1|     M|        ENG|              13|
|    1|     F|        ENG|              40|
|    1|     M|        USA|               1|
|    1|     F|        USA|               4|
|    1|      |        USA|               3|
|    2|     M|        FRA|               4|
|    2|     F|        FRA|               1|
|    2|     M|        USA|              10|
|    2|     F|        USA|               4|
|    2|      |        USA|              60|
+-----+------+-----------+----------------+

你可以使用我上面解释的逻辑

import org.apache.spark.sql.functions._
df.withColumn("Gender", when($"NumberOfCustomer" < 10, lit("Other")).otherwise($"Gender"))
  .withColumn("Nationality", when($"NumberOfCustomer" < 10, lit("Other")).otherwise($"Nationality"))
  .groupBy("Month","Gender","Nationality").agg(sum("NumberOfCustomer").as("NumberOfCustomer"))
  .show()

你应该得到你想要的结果

【讨论】:

  • 这对我帮助很大,但为了获得正确的结果,我必须先 groupBy,然后在每个 WithColumn 之后聚合。感谢您的帮助!
  • 我在回答中也使用了 groupBy 和聚合 :)
【解决方案2】:

您可以将数据转换为DataFrame

val df = rdd.toDF.select(
  $"_1._1" as "month", $"_1._2" as "gender", $"_1._3" as "nationality", 
  $"_2" as "number_of_customers"
).na.replace(Seq("gender", "nationality"), Map("" -> "other"))

并使用过滤器应用多维数据集:

val aggregates = df
  .cube($"month", $"gender", $"nationality")
  .agg(sum("number_of_customers") as "number_of_customers")
  .where($"number_of_customers" >= 10)

这会生成多个集合:

+-----+------+-----------+-------------------+
|month|gender|nationality|number_of_customers|
+-----+------+-----------+-------------------+
| null|  null|        USA|                 20|
|    1| other|       null|                 15|
|    1| other|        USA|                 13|
|    1|  null|       null|                 41|
| null| other|        USA|                 13|
|    1|     M|       null|                 18|
| null|     M|       null|                 18|
| null|  null|        FRA|                 12|
| null|  null|       null|                 41|
| null| other|       null|                 15|
|    1|  null|        FRA|                 12|
|    1|  null|        USA|                 20|
+-----+------+-----------+-------------------+

您可以稍后过滤此结果,例如查找最大、不重叠的关卡集,并且聚合结果应该足够小,以便使用本地集合上的标准集操作来完成。

例如:

val sets = df
   .select(df.columns.map(c => concat_ws("_", lit(c), col(c))): _*)
    .collect.map(row => row.toSeq.collect { case s: String => s }.toSet )

sets.filter(s => !sets.contains(s subsetOf _)).foreach(println)
// Set(month_1, gender_M, nationality_FRA, number_of_customers_8)
// Set(month_1, gender_F, nationality_FRA, number_of_customers_2)
// Set(month_1, gender_other, nationality_FRA, number_of_customers_2)
// Set(month_1, gender_M, nationality_other, number_of_customers_7)
// Set(month_1, gender_F, nationality_other, number_of_customers_2)
// Set(month_1, gender_M, nationality_USA, number_of_customers_3)
// Set(month_1, gender_F, nationality_USA, number_of_customers_4)
// Set(month_1, gender_other, nationality_USA, number_of_customers_13)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-04-07
    • 1970-01-01
    • 1970-01-01
    • 2022-01-16
    • 2019-07-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多