【问题标题】:Apache Spark's RDD splitting according to the particular sizeApache Spark 的 RDD 根据特定大小进行拆分
【发布时间】:2016-03-03 02:21:19
【问题描述】:

我正在尝试从文本文件中读取字符串,但我想根据特定大小限制每一行。例如;

这是我代表的文件。

aaaaa\nbbb\nccccc

当试图通过sc.textFile读取这个文件时,RDD会出现这个。

scala> val rdd = sc.textFile("textFile")
scala> rdd.collect
res1: Array[String] = Array(aaaaa, bbb, ccccc)

但我想限制这个 RDD 的大小。例如,如果限制是3,那么我应该得到这样的。

Array[String] = Array(aaa, aab, bbc, ccc, c)

做到这一点的最佳性能方式是什么?

【问题讨论】:

  • 所以你想忽略行边界并分成n字符组?几乎可以肯定,在 Spark 外部将其预处理为长度为 n 的行,然后使用 textFile 读取它会更快

标签: scala apache-spark rdd


【解决方案1】:

不是一个特别有效的解决方案(也不可怕),但你可以这样做:

val pairs = rdd
  .flatMap(x => x)  // Flatten
  .zipWithIndex  // Add indices
  .keyBy(_._2 / 3)  // Key by index / n

// We'll use a range partitioner to minimize the shuffle 
val partitioner = new RangePartitioner(pairs.partitions.size, pairs)

pairs
  .groupByKey(partitioner)  // group
  // Sort, drop index, concat
  .mapValues(_.toSeq.sortBy(_._2).map(_._1).mkString("")) 
  .sortByKey()
  .values

可以通过传递显式填充分区所需的数据来避免洗牌,但编码需要一些努力。请参阅我对Partition RDD into tuples of length n 的回复。

如果您可以在分区边界上接受一些未对齐的记录,那么简单的 mapPartitions 与分组应该以更低的成本来解决问题:

rdd.mapPartitions(_.flatMap(x => x).grouped(3).map(_.mkString("")))

也可以使用滑动RDD:

rdd.flatMap(x => x).sliding(3, 3).map(_.mkString(""))

【讨论】:

    【解决方案2】:

    无论如何,您都需要读取所有数据。除了映射每条线并对其进行修剪之外,您无能为力。

    rdd.map(line => line.take(3)).collect()
    

    【讨论】:

    • 我不认为是这样。再看看预期的输出。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-12-07
    • 2017-02-17
    • 1970-01-01
    • 2017-02-24
    • 2014-05-13
    • 2016-08-01
    • 2014-09-22
    相关资源
    最近更新 更多