【发布时间】: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 毫秒的数据
对于每个文件,我想将小时数:分钟添加到秒,以便我的时间看起来像小时-分钟-秒。但问题是我不希望这个过程顺序发生
我正在使用 for-each 函数来获取每个文件名 -> 然后在文件中创建数据的 RDD 并添加上面指定的时间。
但我想要的是 每个文件 应该转到一个节点来处理和计算时间,而不是 一个文件中的数据 跨节点分发计算时间。
谢谢。
问候,
维奈·乔格勒卡
【问题讨论】:
-
我想不出与函数
wholeTextFiles有什么不同。我的意思是,理论上你可以创建一个基于文件名的分区器,但据我了解,这类似于wholeTextFiles所做的。 -
但 wholeTextFiles 不保留文件中的列。它将创建一个包含我不想要的所有数据的字符串。
-
那么,你想解析 .csv 文件。为此,您可以使用 github.com/databricks/spark-csv 。你试过了吗?
-
解析部分在后我想对数据进行分区,使得一个分区包含一个文件
-
你为什么要这个?您能否解释一下您试图以这种方式解决的问题。我有一种感觉,你在这里走错了路。
标签: scala csv apache-spark bigdata