【问题标题】:How to split the input file in Apache Spark如何在 Apache Spark 中拆分输入文件
【发布时间】:2014-12-24 17:37:54
【问题描述】:

假设我有一个大小为 100MB 的输入文件。它包含大量 CSV 格式的点(经纬对)。我应该怎么做才能在 Apache Spark 中将输入文件拆分为 10 个 10MB 的文件,或者如何自定义拆分。

注意:我想处理每个映射器中点的子集。

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    Spark 的抽象不提供明确的数据拆分。但是,您可以通过多种方式控制并行度。

    假设您使用 YARN,HDFS 文件会自动拆分为 HDFS 块,并在 Spark 操作运行时同时处理它们。

    除了 HDFS 并行性,考虑使用带有 PairRDD 的分区器。 PairRDD 是键值对 RDD 的数据类型,分区器管理从键到分区的映射。默认分区程序读取spark.default.parallelism。分区器有助于控制数据的分布及其在 PairRDD 特定操作中的位置,例如,reduceByKey

    查看以下有关 Spark 数据并行性的文档。

    http://spark.apache.org/docs/1.2.0/tuning.html

    【讨论】:

    • 场景:我有 50 个点和一个目标点 p0。我必须在这 50 个点中找到最接近 p0 的点。所以我决定将 50 个点分成 5 个 10 个点,并在每 10 个点上并行运行最近邻算法。我从我的一点 MapReduce 知识中知道,每个映射器都只需要一条线而无需自定义。因此,每个映射器都有一个点进行操作,而不是 10 个点。如何解决这个问题?
    • Suztomo - textFile() 上还有一个 minPartitions 参数,可以控制将文件加载到多少个分区。 @Chandan - 如果您的 RDD 没有足够的分区,请在运行计算之前尝试使用 RDD.repartition(N) 显式重新分区。更多,更小的分区将为每个任务(我不认为我们在 Spark 中谈论“映射器”)减少工作量。
    • JavaRDD<String> lines = ctx.textFile("/home/hduser/Spark_programs/file.txt").cache(); lines.repartition(2); List<Partition> partitions = lines.partitions(); System.out.println(partitions.size()); 只给一个分区。
    【解决方案2】:

    通过 Spark API 搜索后,我发现了一种方法 partition,它返回 JavaRDD 的分区数。在创建 JavaRDD 时,我们已将其重新分区为@Nick Chammas 告知的所需分区数。

    JavaRDD<String> lines = ctx.textFile("/home/hduser/Spark_programs/file.txt").repartition(5);
    List<Partition> partitions = lines.partitions();
    System.out.println(partitions.size());
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-03-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多