【问题标题】:Java spark Enumeration in parallelJava spark 并行枚举
【发布时间】:2019-03-07 14:39:20
【问题描述】:

我在 Java Spark 中运行了以下代码:

ZipFile zipFile = new ZipFile(zipFilePath);
Enumeration<? extends ZipEnter> entries = zipFiles.entries();
while(entries.hasMoreElements()) {
    ZipEntry entry = entries.nextElement();
    //my logic...
}

我想用 Spark 或 Java 并行执行上面的代码,我该怎么做?

谢谢

【问题讨论】:

    标签: java apache-spark parallel-processing


    【解决方案1】:

    下面的代码将分别在java和scala中同时处理枚举中每个条目的逻辑。

    在 Java 中

    entriesList = Collections.list(enumeration);
    List<CompletableFuture<ZipEnter>> futureList = entriesList.stream().(x -> CompletableFuture. supplyAsync(() -> {
        //logic
    }).collect(Collectors.toList());
    CompletableFuture.allof(futureList);
    

    在 Scala 中

        entriesList = // to scala list
    
        Future[ZipEnter] futureList = entriesList.map(x => Future{
            // logic
        })
    
        Future.sequence(futureList)
    

    希望对你有帮助。

    【讨论】:

    • 我应该在 entryList.stream() 旁边写什么。 ? entryList.stream().allMatch?
    • 没有。你根本不需要 allMatch 。您只需将 //logic 替换为任何逻辑即可。我只是想将ZipEnter 列表转换为CompletableFuture&lt; ZipEnter&gt; 列表。
    【解决方案2】:
    import org.apache.spark.SparkConf;
    import org.apache.spark.api.java.JavaSparkContext;
    import org.apache.spark.api.java.function.VoidFunction;
    
    import java.io.BufferedInputStream;
    import java.io.DataInputStream;
    import java.io.File;
    import java.io.FileInputStream;
    import java.util.Arrays;
    import java.util.List;
    import java.util.Objects;
    
    public class ParallelEnumeration {
        public static void main(String[] args) {
            String zipFilePath = "/ZipDir/";
            File zipFiles = new File(zipFilePath);
            final List<File> files = Arrays.asList(Objects.requireNonNull(zipFiles.listFiles()));
            // configure spark
            SparkConf sparkConf = new SparkConf().setAppName("Print Elements of RDD")
                    .setMaster("local[*]");
            // start a spark context
            JavaSparkContext jsc = new JavaSparkContext(sparkConf);
    
            // parallelize the file collection to two partitions
            jsc.parallelize(files, 2)
                    .filter(file -> { // This filter is optional if the directory contains only zip files
                        // https://stackoverflow.com/questions/33934178/how-to-identify-a-zip-file-in-java
                        DataInputStream in = new DataInputStream(new BufferedInputStream(new FileInputStream(file)));
                        int test = in.readInt();
                        in.close();
                        return test == 0x504b0304;
                    }).foreach((VoidFunction<File>) file -> System.out.println(file.getName()));
    
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-12-10
      • 2014-02-27
      • 2014-08-19
      • 1970-01-01
      相关资源
      最近更新 更多