在 Apache Spark 中对 RDD 进行分区,使得一个分区包含在一个文件中

Vin*_*kar 5 csv scala bigdata apache-spark

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

val combineRDD = sc.textFile("D://release//CSVFilesParellel//*.csv")
Run Code Online (Sandbox Code Playgroud)

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

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

我如何实现这一目标?

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

00-00

时间(以秒为单位) Measure1 Measure2 Measure3..... Measuren

0

0.25

0.50

0.75

1

...

3600


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

2.第一列是第二列,分为4个部分,每部分250 ms,数据记录250 ms

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

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

  3. 但我想要的是每个文件都应该去一个节点来处理和计算时间,而不是一个文件中的数据分布在节点之间来计算时间。

谢谢你。

问候,

维奈·乔格卡

Kra*_*tam 1

让我们回到基础知识。

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

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

要在 Spark 中使用自定义输入格式,请遵循

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