【发布时间】:2014-12-24 17:37:54
【问题描述】:
假设我有一个大小为 100MB 的输入文件。它包含大量 CSV 格式的点(经纬对)。我应该怎么做才能在 Apache Spark 中将输入文件拆分为 10 个 10MB 的文件,或者如何自定义拆分。
注意:我想处理每个映射器中点的子集。
【问题讨论】:
标签: apache-spark
假设我有一个大小为 100MB 的输入文件。它包含大量 CSV 格式的点(经纬对)。我应该怎么做才能在 Apache Spark 中将输入文件拆分为 10 个 10MB 的文件,或者如何自定义拆分。
注意:我想处理每个映射器中点的子集。
【问题讨论】:
标签: apache-spark
Spark 的抽象不提供明确的数据拆分。但是,您可以通过多种方式控制并行度。
假设您使用 YARN,HDFS 文件会自动拆分为 HDFS 块,并在 Spark 操作运行时同时处理它们。
除了 HDFS 并行性,考虑使用带有 PairRDD 的分区器。 PairRDD 是键值对 RDD 的数据类型,分区器管理从键到分区的映射。默认分区程序读取spark.default.parallelism。分区器有助于控制数据的分布及其在 PairRDD 特定操作中的位置,例如,reduceByKey。
查看以下有关 Spark 数据并行性的文档。
【讨论】:
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()); 只给一个分区。
通过 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());
【讨论】: