【问题标题】:Spark use reduceByKey on a nested Key-Value structureSpark 在嵌套的键值结构上使用 reduceByKey
【发布时间】:2017-10-07 12:57:54
【问题描述】:

我的数据如下:

客户1|项目1:x1,x2,x3;项目2:x1,x4,x5; item1:x1,x3,x6|time1|url
客户1|项目1:x1,x7,x3;项目2:x1,x4,x5; item3:x5|time2|url2
客户2|项目1:x1,x7,x3; item3:x5|time3|url3

我想 ReduceByKey 相同的 customerIds 和 mapValues 以获得每个 customerId 的不同项目的联合:

客户1|项目1:x1,x2,x3;项目2:x1,x4,x5;项目1:x1,x3,x6;项目1:x1,x7,x3;项目3:x5

我可以通过以下方式实现:

val line = spark.sparkContext.textFile(args(0))
val record = line.map(l=>l.split("\|")).map(l=>(l(0),l(1))).reduceByKey((x,y) => x. union(y)).mapValues(x=>x.distinct)

现在,我希望第二列中的每个项目也都是唯一的,并且应该使用 union 和 distinct 连接同一键中的所有值,以获得类似:

客户1|项目1:x1,x2,x3,x6,x7;项目2:x1,x4,x5;项目3:x5

一旦完成,我想选择每个 x 的所有频率,例如:x1:2, x2:1 .... 并用我得到的频率更新了 customerId 的 x(1-10) 向量。

这可以在 spark 中实现吗?

【问题讨论】:

    标签: scala apache-spark spark-dataframe


    【解决方案1】:

    是的,您当然可以在 Spark 中做到这一点!然而,您处理问题的方式看起来有点困难。

    所以我可以展示一个完整的可复制粘贴到 REPL 示例,假设您的数据存储在字符串中(不是 args(0) 文件)

    val data = """Customer1| item1:x1,x2,x3; item2:x1,x4,x5; item1:x1,x3,x6|time1|url
    Customer1| item1:x1,x7,x3; item2:x1,x4,x5; item3:x5|time2|url2
    Customer2| item1:x1,x7,x3; item3:x5|time3|url3"""
    

    你称之为“line”的RDD可以读入RDD“rdd”为

    val rdd = sc.parallelize(data.split("\n"))
    

    到目前为止还没有什么新鲜事。下一步是重要的一步。我们无需逐层进行计数和汇总,而是可以准备数据以一次性完成所有操作。这更易读,也更高效,因为它是一个单一的 map 后跟一个 reduce。

    val mapped= rdd.flatMap(line => {
       val arr = line.split("\\|")
       val customer = arr(0)
       val items = arr(1)
       val time = arr(2)
       val url = arr(3)
    
       items.split(";").flatMap(item => {
          val itemKey = item.split(":")(0)
          val itemValues = item.split(":")(1).split(",")
    
          itemValues.map(value => (customer, itemKey, value, time, url))
       })
    })
    

    我们可以看到里面有什么我们可以用mapped.toDF("customer", "itemId", "itemValue", "time", "url").show很好地打印出来

    +---------+------+---------+-----+----+
    | customer|itemId|itemValue| time| url|
    +---------+------+---------+-----+----+
    |Customer1| item1|       x1|time1| url|
    |Customer1| item1|       x2|time1| url|
    |Customer1| item1|       x3|time1| url|
    |Customer1| item2|       x1|time1| url|
    |Customer1| item2|       x4|time1| url|
    |Customer1| item2|       x5|time1| url|
    |Customer1| item1|       x1|time1| url|
    |Customer1| item1|       x3|time1| url|
    |Customer1| item1|       x6|time1| url|
    |Customer1| item1|       x1|time2|url2|
    |Customer1| item1|       x7|time2|url2|
    |Customer1| item1|       x3|time2|url2|
    |Customer1| item2|       x1|time2|url2|
    |Customer1| item2|       x4|time2|url2|
    |Customer1| item2|       x5|time2|url2|
    |Customer1| item3|       x5|time2|url2|
    |Customer2| item1|       x1|time3|url3|
    |Customer2| item1|       x7|time3|url3|
    |Customer2| item1|       x3|time3|url3|
    |Customer2| item3|       x5|time3|url3|
    +---------+------+---------+-----+----+
    

    最后我们可以计算并归约成你需要的向量:

    val reduced = mapped.map{case (customer, itemKey, itemValue, time, url) => ((customer, itemKey, itemValue), 1)}.
       reduceByKey(_+_).
       map{case ((customer, itemKey, itemValue), count) => (customer, itemKey, itemValue, count)}
    

    并查看它:reduced.toDF("customer", "itemKey", "itemValue", "count").show

    +---------+-------+---------+-----+                                             
    | customer|itemKey|itemValue|count|
    +---------+-------+---------+-----+
    |Customer1|  item1|       x2|    1|
    |Customer1|  item1|       x1|    3|
    |Customer2|  item1|       x7|    1|
    |Customer1|  item1|       x6|    1|
    |Customer1|  item1|       x7|    1|
    |Customer2|  item1|       x3|    1|
    |Customer2|  item3|       x5|    1|
    |Customer1|  item2|       x5|    2|
    |Customer1|  item2|       x4|    2|
    |Customer1|  item2|       x1|    2|
    |Customer1|  item3|       x5|    1|
    |Customer1|  item1|       x3|    3|
    |Customer2|  item1|       x1|    1|
    +---------+-------+---------+-----+
    

    如果您需要将所有内容分组到向量的 Array/Seq 表示中,您可以通过进一步聚合数据来做到这一点。希望这会有所帮助!

    【讨论】:

    • 还有一些值的时间和 url 不存在,在这种情况下 arr(2) 和 arr(3) 将因 ArrayIndexOutOfBoundsException 而失败。是否可以过滤具有 4 列的行。像 line.split("\\|")).filter(l=>l.length == 4) 我可以忽略没有 url 和时间的数据。
    • 如果不需要的话,只需从元组中删除这些列。或者,import scala.util.Try 然后将这些行更新为 val time = Try(Some(arr(2))).getOrElse(None)val url = Try(Some(arr(3))).getOrElse(None)
    • 取决于您是否需要这些行中的值。如果你不这样做,那么你可以按照你的建议过滤。如果你这样做了,那么看看之前的评论:)
    • 那行得通..谢谢! :) 我需要的reducer 是val reduced = mapped.map{case (customer, itemKey, itemValue, time, url) => ((customer, itemValue), 1)}. reduceByKey(_+_). map{case ((customer, itemValue), count) => (customer, itemValue, count)} 关于最后一部分.. 我如何为每个customerId 创建一个itemValue 计数向量。例如。 if header[x0,x1,x2,x3,x4,x5,x6,x7,x8,x9] customerId2: item_freq[0 1 0 0 3 0 2 0 0 0] 根据每个 x 获得的值,如果没有获得则为 0在哪里
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-03-14
    • 1970-01-01
    • 2016-12-27
    • 1970-01-01
    • 2015-07-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多