【问题标题】:Partition RDD in Apache Spark such that one partition consists on one fileApache Spark 中的分区 RDD,使得一个分区包含在一个文件中
【发布时间】:2016-10-08 15:39:30
【问题描述】:

我正在创建一个这样的 2.csv 文件的 RDD

val combineRDD = sc.textFile("D://release//CSVFilesParellel//*.csv")

然后我想在这个 RDD 上定义自定义分区,这样一个分区必须包含一个文件。 这样每个分区,即一个 csv 文件在一个节点上处理,以加快数据处理速度

是否可以根据文件大小或一个文件中的行数或一个文件的文件结尾字符编写自定义分区程序?

我如何做到这一点?

一个文件的结构如下所示:

00-00

时间(秒)Measure1 Measure2 Measure3..... Measuren

0

0.25

0.50

0.75

1

...

3600


1.第一行数据包含小时数:mins 每个文件包含 1 小时或 3600 秒的数据

2.第一列是第二列,分为 4 部分,每部分 250 毫秒,记录 250 毫秒的数据

  1. 对于每个文件,我想将小时数:分钟添加到秒,以便我的时间看起来像小时-分钟-秒。但问题是我不希望这个过程顺序发生

  2. 我正在使用 for-each 函数来获取每个文件名 -> 然后在文件中创建数据的 RDD 并添加上面指定的时间。

  3. 但我想要的是 每个文件 应该转到一个节点来处理和计算时间,而不是 一个文件中的数据 跨节点分发计算时间。

谢谢。

问候,

维奈·乔格勒卡

【问题讨论】:

  • 我想不出与函数wholeTextFiles 有什么不同。我的意思是,理论上你可以创建一个基于文件名的分区器,但据我了解,这类似于 wholeTextFiles 所做的。
  • 但 wholeTextFiles 不保留文件中的列。它将创建一个包含我不想要的所有数据的字符串。
  • 那么,你想解析 .csv 文件。为此,您可以使用 github.com/databricks/spark-csv 。你试过了吗?
  • 解析部分在后我想对数据进行分区,使得一个分区包含一个文件
  • 你为什么要这个?您能否解释一下您试图以这种方式解决的问题。我有一种感觉,你在这里走错了路。

标签: scala csv apache-spark bigdata


【解决方案1】:

简单的答案,无需质疑您为什么这样做。分别加载文件,以便您知道正在加载的文件名

// create firstRDD containing a new attribute `filename=first.csv`
val firstRDD = sc.textFile("D://release//CSVFilesParellel//first.csv")
    .map(line => new CsvRecord(line))

// create secondRDD containing a new attribute `filename=second.csv`
val secondRDD = sc.textFile("D://release//CSVFilesParellel//second.csv")
    .map(line => new CsvRecord(line))

// now create a pair RDD and re-partition on the filename
val partitionRDD = firstRDD.union(secondRDD)
    .map(csvRecord => (csvRecord.filename,csvRecord))
    .partitionBy(customFilenamePartitioner)

以下引用来自here

要实现您的 customFilenamePartitioner,您需要将 org.apache.spark.Partitioner 类并实现三种类型 方法:

NumPartitions : Int,返回您将根据要求创建的分区数。

getPartition(key: Any) : Int,返回分区ID范围 对于给定的键,从 (0 到 numPartitions-1)。

equals(), :标准的 Java 编程相等方法。这是 实施很重要,因为 Spark 应用程序需要测试 您的 Partitioner 对象按其自己的条款针对其他实例 它决定您的两个 RDD 是否以与其相同的方式进行分区 是必需的。

请记住,重新分区很可能会触发昂贵的 shuffle,因此除非您要重复查询这个新分区的 RDD,否则最好以另一种方式解决您的问题。

【讨论】:

    【解决方案2】:

    让我们回到基础。

    1. BigData 的哲学将流程转移到数据而不是数据转移到处理。这样可以提高并行度,从而提高 I/O 吞吐量
    2. 一个分区器占用一个文件将减少并行度而不是增加。
    3. 实现此目的的最简单方法是使用 textInpuTFormat 并通过 gzip 或 lzo 压缩您的输入文件(不应该进行 lzo 索引)。
    4. Gzip 不可拆分将强制一个文件转到一个分区,但这绝不会有助于提高任何类型的吞吐量

    5. 编写自定义输入格式从 FileInputFormat 扩展并提供您的 splitlogic 和 recordReader 逻辑。

    要在 spark 中使用自定义输入格式,请关注

    http://bytepadding.com/big-data/spark/combineparquetfileinputformat/

    【讨论】:

      猜你喜欢
      • 2016-08-01
      • 2015-06-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-09-18
      • 2016-01-04
      • 2010-11-28
      • 1970-01-01
      相关资源
      最近更新 更多