【问题标题】:Spark RDD: partitioning according to text file formatSpark RDD:根据文本文件格式进行分区
【发布时间】:2018-12-07 22:39:43
【问题描述】:

我有一个包含数十 GB 数据的文本文件,我需要从 HDFS 加载并并行化为 RDD。此文本文件使用以下格式描述项目。请注意,字母字符串不存在(每行的含义是隐含的),并且每行可以包含空格以分隔不同的值:

0001  (id)
1000 1000 2000 (dimensions)
0100           (weight)
0030           (amount)
0002  (id)
1110 1000 5000 (dimensions)
0220           (weight)
3030           (amount)

我认为并行化此文件的最直接方法是将其从本地文件系统上传到 HDFS,然后通过执行 sc.textFile(filepath) 创建 RDD。但是,在这种情况下,分区将取决于文件对应的 HDFS 拆分。

上述方法的问题是每个分区可能包含不完整的项目。例如:

分区 1

0001           (id)
1000 1000 2000 (dimensions)
0100           (weight)
0030           (amount)
0002           (id)
1110 1000 5000 (dimensions)

分区 2

0220           (weight)
3030           (amount)

因此,当我们为每个分区调用一个方法并将其对应的数据块传递给它时,它将收到标识为 0002 的项目的不完整规范。这将导致在调用内部执行的计算输出错误方法。

为了避免这个问题,对这个 RDD 进行分区或重新分区的最有效方法是什么?可以指定每个分区的行数为4的倍数吗?如果是,应该由Hadoop还是Spark来完成?

【问题讨论】:

    标签: apache-spark hadoop rdd hadoop-partitioning


    【解决方案1】:

    加载文本文件获取RDD[String],然后使用zipWithIndex转换为RDD[(String, Long)],其中元组中的第二个属性是RDD中元素的索引号。

    用它的元素索引压缩这个 RDD。排序首先基于分区索引,然后是每个分区内项目的排序。所以第一个分区的第一项获得索引0,最后一个分区的最后一项获得最大的索引。

    • 使用索引作为行号(从 0 开始),我们可以对属于记录的行进行分组。例如。 [0, 1, 2, 3, 4, 5, 6, 7, 8, 9, ...
    • 由于我们知道每条记录(确切地说)跨越 4 行,索引除以 4 的整数除法(我们称之为 idx_div)。这将导致前四行有 0 作为idx_div,接下来的四行将得到 1 作为idx_div,依此类推。例如。 [0, 0, 0, 0, 1, 1, 1, 1, 2, 2, ...。这可用于对属于一条记录的所有(四)行进行分组,以便进一步解析和处理


    case class Record(id:String, dimensions:String, weight:String, amount:String)
    val lines = sc.textFile("...")
    val records = lines
        .zipWithIndex
        .groupBy(line_with_idx => (line_with_idx._2 / 4))  // groupBy idx_div
        .map(grouped_record => {
            val (idx_div:Long, lines_with_idx:Iterable[(String, Long)]) = grouped_record
            val lines_with_idx_list = lines_with_idx.toList.sortBy(_._2)  // Additional check to ensure ordering
            val lines_list = lines_with_idx_list.map(_._1)
            val List(id:String, dimensions:String, weight:String, amount:String) = lines_list
            new Record(id, dimensions, weight, amount)
        })
    

    【讨论】:

    • 很好的答案,尽管如果您添加有关代码及其含义的更详细说明,这将非常有用。谢谢
    【解决方案2】:

    为什么不简单地在将文件放入 HDFS 之前对行进行分组以避免这个问题?

    xargs -L4 echo < file
    hdfs dfs -put file /your/path
    

    你的数据看起来像

    0001  1000  0100  0030 
    0002  1110  0220  3030
    

    如果这样做,您可以使用更优化的 Spark DataFrames API 读取数据 比 RDD 更丰富,并为您提供更丰富的 API 和性能来编写您的应用程序。

    【讨论】:

    • 不幸的是,我忘了解释每个原始行可以包含空格来分隔不同的值。我刚刚编辑了这个问题以包括这个方面。如您所示对行进行分组会与此冲突,因为空格不能替换换行符作为每个原始行的分隔符。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-04-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多