【问题标题】:Accessing a global lookup Apache Spark访问全局查找 Apache Spark
【发布时间】:2015-03-30 10:26:09
【问题描述】:

我有一个 csv 文件列表,每个文件都有一堆类别名称作为标题列。每行都是一个具有布尔值 (0, 1) 的用户列表,无论他们是否属于该类别。每个 csv 文件没有相同的标题类别集。

我想在所有具有以下输出的文件中创建一个复合 csv:

  1. 标题是所有标题的联合
  2. 每一行都是一个唯一的用户,具有对应于类别列的布尔值

我想解决这个问题的方法是为每个带有“1”的单元格创建一个 user_id 和一个唯一 category_id 的元组。然后为每个用户减少所有这些列以获得最终输出。

如何开始创建元组?我可以对所有类别进行全局查找吗?

示例数据:

File 1
user_id,cat1,cat2,cat3
21321,,,1,
21322,1,1,1,
21323,1,,,

文件 2

user_id,cat4,cat5
21321,1,,,
21323,,1,,

输出

user_id,cat1,cat2,cat3,cat4,cat5
21321,,1,1,,,
21322,1,1,1,,,
21323,1,1,,,,

【问题讨论】:

  • 示例数据有助于说明您的输入和期望的结果。
  • 遇到冲突怎么办?例如一个 csv 表示 user_id, cat2 = 0 和另一个 user_id, cat2 = 1 ?另外,到目前为止,您尝试过什么?
  • 值的存在优先。到目前为止,我一直在试图说服它对于 map reduce 功能来说是一种糟糕的格式。我们控制数据的格式,因此计划对其进行更改。但我很好奇是否有其他人遇到过这类问题。

标签: apache-spark bigdata


【解决方案1】:

如果您只需要结果值,则无需将此作为两步过程。 一个可能的设计: 1/解析你的csv。你没有提到你的数据是否在分布式 FS 上,所以我假设它不是。 2/ 将您的 (K,V) 对输入到可变并行化(以利用 Spark)映射中。 伪代码:

val directory = ..
mutable.ParHashMap map = new mutable.ParHashMap()
while (files[i] != null)
{
  val file = directory.spark.textFile("/myfile...")
  val cols = file.map(_.split(","))

  map.put(col[0], col[i++])
}

然后您可以通过地图上的迭代器访问您的 (K/V) 元组。

【讨论】:

  • 我添加了一个示例,进一步澄清了这个问题。这不会解决它..
  • 示例令人困惑:我认为可能的单元格值在 (0,1) 中,但似乎也有 '21322' 值?
  • 如果不清楚的话,我很糟糕。数字是用户 ID。其他所有内容都是 1 或 '' 对应于该值是否存在。
【解决方案2】:

从某种意义上说,问题的标题可能具有误导性,因为不需要全局查找来解决手头的问题。

在大数据中,有一个指导大多数解决方案的基本原则:分而治之。在这种情况下,输入的 CSV 文件可以分为 (user,category) 的元组。 包含任意数量类别的任意数量的 CSV 文件都可以转换为这种简单的格式。上一步并集得到的 CSV 结果,提取当前类别的总 nr 并进行一些数据转换以得到所需的格式。

在代码中,该算法如下所示:

import org.apache.spark.SparkContext._

val file1 = """user_id,cat1,cat2,cat3|21321,,,1|21322,1,1,1|21323,1,,""".split("\\|")
val file2 = """user_id,cat4,cat5|21321,1,|21323,,1""".split("\\|")
val csv1 = sparkContext.parallelize(file1)
val csv2 = sparkContext.parallelize(file2)

import org.apache.spark.rdd.RDD
def toTuples(csv:RDD[String]):RDD[(String, String)] = {
  val headerLine = csv.first
  val header = headerLine.split(",")
  val data = csv.filter(_ != headerLine).map(line => line.split(","))
  data.flatMap{elem => 
    val merged = elem.zip(header)
    val id = elem.head
    merged.tail.collect{case (v,cat) if v == "1" => (id, cat)}
  }
}

val data1 = toTuples(csv1)
val data2 = toTuples(csv2)
val union = data1.union(data2)
val categories = union.map{case (id, cat) => cat}.distinct.collect.sorted //sorted category names
val categoriesByUser = union.groupByKey.mapValues(v=>v.toSet)
val numericCategoriesByUser = categoriesByUser.mapValues{catSet => categories.map(cat=> if (catSet(cat)) "1" else "")}
  val asCsv = numericCategoriesByUser.collect.map{case (id, cats)=> id + "," + cats.mkString(",")}

结果:

21321,,,1,1,
21322,1,1,1,,
21323,1,,,,1

(生成标题很简单,留给读者练习)

【讨论】:

    猜你喜欢
    • 2016-12-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-06
    • 2012-10-11
    • 1970-01-01
    • 1970-01-01
    • 2015-05-10
    相关资源
    最近更新 更多