【发布时间】: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