【问题标题】:Modified Structure of RDD in SparkSpark中RDD的修改结构
【发布时间】:2017-02-04 03:59:47
【问题描述】:

我是 spark/scala 的新手。

val First: RDD[((Short, String), (Int, Double, Int))]

这是RDD的结构。我想修改这个结构,如下所示:

val First: RDD[(Short, String , Int, Double, Int)]

因为我有另一个结构不同的 RDD,我想联合这两个 RDD。 (在 UNION 操作中结构必须相同)。

请给我建议一个选项。

【问题讨论】:

  • 无汗:First.map { case ((x,y),(z,w)) => (x,y,z,w) }
  • @Alec 我试过这个,但由于数据量很大,所以这会降低性能。因为 Map 会一一迭代数据。
  • 请给我一些解决方案,我可以在不迭代数据的情况下更改结构
  • 其实不会。像map 这样的转换在 Spark 中是惰性执行的。 map 最终与完成转换链的任何操作同时计算 - 没有中间步骤。无论如何,减速都会在您的集群上并行化,所以如果您的集群甚至无法处理这个问题,它可能无法处理您之后计划做的任何其他事情......
  • @Alec 我的更新答案非常同意你的观点(对吗?:)),所以 Darshan 不要感到困惑,我真的同意 Alec!

标签: scala function apache-spark distributed-computing bigdata


【解决方案1】:

只需映射您的数据,如下所示:

First.map{ case ( (x, y), (k, z, w) ) => (x, y, k, z, w) }

为了写这个map函数,你必须检查你的RDD的格式,((Short, String), (Int, Double, Int)),也就是我写的(x, y), (k, z, w),然后在@987654325的右边写上你想要的格式@。


编辑评论:

因为 Map 会一个一个地迭代数据

仅在动作发生时应用转换,因此map() 以分布式方式工作得非常好。每个分区都会在其数据中应用映射函数。

虽然这不是一个成本很高的操作,所以不要专注于那个,专注于你的加入,这是一个繁重的操作。如果您的集群中有相应的资源,那么对于您的数据量,map 函数应该是便宜的。

【讨论】:

  • 你有没有其他选项可以修改结构而不进行迭代(因为 Map 会一一迭代数据)
猜你喜欢
  • 1970-01-01
  • 2014-06-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多