【问题标题】:Processing several files one-by-one separately by SparkSpark分别处理多个文件
【发布时间】:2017-07-11 12:35:45
【问题描述】:

我在使用 Apache Spark 实施一个工作流方面需要帮助。我的任务在下一个:

  1. 我有几个 CSV 文件作为源数据。注意:这些文件可能有不同的布局
  2. 我有元数据,其中包含我需要如何解析每个文件的信息(这不是问题)
  3. 主要目标:结果是带有几个附加列的源文件。我必须更新每个源文件而不加入一个输出范围。例如:源 10 个文件 -> 10 个结果文件,每个结果文件只有对应源文件的数据。

据我所知,Spark 可以通过掩码打开许多文件:

var source = sc.textFile("/source/data*.gz");

但在这种情况下,我无法识别文件的哪一行。如果我得到源文件列表并尝试通过以下场景进行处理:

JavaSparkContext sc = new JavaSparkContext(...);
List<String> files = new ArrayList() //list of source files full name's
for(String f : files)
{
   JavaRDD<String> data = sc.textFile(f);
   //process this file with Spark
   outRdd.coalesce(1, true).saveAsTextFile(f + "_out"); 
}

但在这种情况下,我将以顺序模式处理所有文件。

接下来是我的问题:如何以并行模式处理多个文件?例如:一个文件 - 一个执行者?

非常感谢您的帮助!

【问题讨论】:

    标签: apache-spark parallel-processing


    【解决方案1】:

    步骤如下

    1. 使用 sparkcontext.wholeTextFiles("/path/to/folder/ contains/all/files")
    2. 上面返回一个RDD,其中key是文件的路径,value是文件的内容
    3. rdd.map(lambda x:x[1]) - 这给你一个只有文件内容的rdd
    4. rdd.map(lambda x: customeFunctionToProcessFileContent(x))
    5. 由于 map 函数是并行工作的,因此您执行的任何操作都会更快且不连续 - 只要您的任务不相互依赖,这是并行性的主要标准

    上述方法适用于默认分区。所以你可能不会得到输入文件数等于输出文件数(因为输出是分区数)。

    您可以根据计数或基于您的数据的任何其他唯一值对 RDD 重新分区,因此最终输出文件计数等于输入计数。这种方法仅具有并行性,但不会达到最佳分区数的性能

    【讨论】:

    • 您好 Ramzy,感谢您的回答,但我还有一个疑问。方法sparkcontext.wholeTextFiles("/path/to/folder/containing/all/files") 打开并读取内存中的文件。据我所知,大多数源文件将有大约 1-3 百万行,但有几个文件的大小可达 2-3 GB。这将在没有任何内存错误的情况下工作?
    • 当你使用sc.textFile或sc.wholeTextFiles时,计算还没有开始。只有当您执行任何操作时才会开始处理,并且基于默认分区来划分数据集。您可以通过 yourRDD.partitions.length 获取分区数并根据需要进行自定义,还可以通过 yourRDD.count() 获取实际 RDD 大小。
    • @Ramzy,wholeTextFiles 将使用路径键和文件整个上下文的值创建 RDD。如果某些文件是 2-3GB 显然会有问题(取决于执行器内存,但无论如何 1 个分区的 GB 太多了)
    【解决方案2】:

    您可以打开常规 java 固定大小的线程池(例如 10 个线程)并从 Callable/Runnable 提交您的 saveAsTextFile 的 spark 作业。 这将提交 10 个并行作业,如果您的 spark 集群中有足够的资源 - 它们将并行执行。 类似以下的东西

    import java.util.ArrayList;
    import java.util.List;
    import java.util.concurrent.Executor;
    import java.util.concurrent.ExecutorService;
    import java.util.concurrent.Executors;
    import java.util.concurrent.Future;
    
    import org.apache.spark.api.java.JavaRDD;
    import org.apache.spark.api.java.JavaSparkContext;
    
    import com.google.common.collect.Lists;
    
    public class Test {
    
        public static void main(String[] argv) {
            final JavaSparkContext sc = new JavaSparkContext(...);
            List<String> files = new ArrayList<>(); //list of source files full name's
            ExecutorService pool = Executors.newFixedThreadPool(10);
            List<Future<?>> futures = new ArrayList<>();
            for(final String f : files)
            {
                Future<?> fut = pool.submit(new Runnable() {
    
                    @Override
                    public void run() {
                        JavaRDD<String> data = sc.textFile(f);
                        //process this file with Spark
                        outRdd.coalesce(1, true).saveAsTextFile(f + "_out"); 
    
                    }
                });
                futures.add(fut);
    
            }
            //waiting for all tasks
            for (Future<?> fut : futures) {
                fut.get();
            }
        }
    }
    

    【讨论】:

    • 谢谢,我认为这是有道理的。我会尝试这种方法。
    • 我想知道线程的任务是如何定义的,它们是如何收集和呈现的。使用这种方法,是否可以实现 10 的并行度? Mapreduce 和 spark 应用程序用于并行处理。请重新审视基础知识,看看它们是否符合要求
    • @Yustas,我添加了一些将您的任务包装在 Runnable 中的代码
    • @Ramzy,看看并尝试自己。这种方法有效。如果您从驱动程序中的不同线程定义火花动作 - 所有这些都将转换为单独的并行作业。 Parallelilsm 将是 10 * 每个文件中的分区数。
    • 绝对可行。但是在线程的情况下,如何设置要处理的文件的限制,然后从中获取结果?如果使用得当,所有这些东西都由 spark/mapreduce 处理。如果线程的使用符合您的要求,欢迎您继续。我只是想了解这个过程。谢谢
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-07-30
    • 2020-06-15
    • 2015-07-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多