是的,您当然可以在 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 表示中,您可以通过进一步聚合数据来做到这一点。希望这会有所帮助!