【问题标题】:Use SparkContext hadoop configuration within RDD methods/closures, like foreachPartition在 RDD 方法/闭包中使用 SparkContext hadoop 配置,例如 foreachPartition
【发布时间】:2016-11-08 12:18:01
【问题描述】:

我正在使用 Spark 读取一堆文件,对它们进行详细说明,然后将它们全部保存为序列文件。我想要的是每个分区有 1 个序列文件,所以我这样做了:

SparkConf sparkConf = new SparkConf().setAppName("writingHDFS")
                .setMaster("local[2]")
                .set("spark.streaming.stopGracefullyOnShutdown", "true");
        final JavaSparkContext jsc = new JavaSparkContext(sparkConf);
        jsc.hadoopConfiguration().addResource(hdfsConfPath + "hdfs-site.xml");
        jsc.hadoopConfiguration().addResource(hdfsConfPath + "core-site.xml");
        //JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(5*1000));

        JavaPairRDD<String, PortableDataStream> imageByteRDD = jsc.binaryFiles(sourcePath);
        if(!imageByteRDD.isEmpty())
            imageByteRDD.foreachPartition(new VoidFunction<Iterator<Tuple2<String,PortableDataStream>>>() {

                @Override
                public void call(Iterator<Tuple2<String, PortableDataStream>> arg0){
                        throws Exception {
                  [°°°SOME STUFF°°°]
                  SequenceFile.Writer writer = SequenceFile.createWriter(
                                     jsc.hadoopConfiguration(), 
//here lies the problem: how to pass the hadoopConfiguration I have put inside the Spark Context? 
Previously, I created a Configuration for each partition, and it works, but I'm sure there is a much more "sparky way"

有人知道如何在 RDD 闭包内部使用 Hadoop 配置对象吗?

【问题讨论】:

    标签: java hadoop apache-spark rdd


    【解决方案1】:

    这里的问题是 Hadoop 配置没有标记为 Serializable,因此 Spark 不会将它们拉入 RDD。它们被标记为Writable,因此 Hadoop 的序列化机制可以编组和解组它们,但 Spark 不能直接使用它

    两个长期修复选项是

    1. 添加对 Spark 中可写序列化的支持。也许SPARK-2421
    2. 使 Hadoop 配置可序列化。
    3. 添加对序列化 Hadoop 配置的显式支持。

    对于使 Hadoop conf 可序列化,您不会遇到任何主要的反对意见;前提是您实现了委托给可写 IO 调用的自定义 ser/deser 方法(并且仅遍历所有键/值对)。我是作为 Hadoop 提交者这么说的。

    更新:下面是创建可序列化类的代码,该类对 Hadoop 配置的内容进行编组。使用val ser = new ConfSerDeser(hadoopConf) 创建它;在您的 RDD 中将其称为 ser.get()

    /*
     * Licensed to the Apache Software Foundation (ASF) under one or more
     * contributor license agreements.  See the NOTICE file distributed with
     * this work for additional information regarding copyright ownership.
     * The ASF licenses this file to You under the Apache License, Version 2.0
     * (the "License"); you may not use this file except in compliance with
     * the License.  You may obtain a copy of the License at
     *
     *    http://www.apache.org/licenses/LICENSE-2.0
     *
     * Unless required by applicable law or agreed to in writing, software
     * distributed under the License is distributed on an "AS IS" BASIS,
     * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
     * See the License for the specific language governing permissions and
     * limitations under the License.
     */
    
     import org.apache.hadoop.conf.Configuration
    
    /**
     * Class to make Hadoop configurations serializable; uses the
     * `Writeable` operations to do this.
     * Note: this only serializes the explicitly set values, not any set
     * in site/default or other XML resources.
     * @param conf
     */
    class ConfigSerDeser(var conf: Configuration) extends Serializable {
    
      def this() {
        this(new Configuration())
      }
    
      def get(): Configuration = conf
    
      private def writeObject (out: java.io.ObjectOutputStream): Unit = {
        conf.write(out)
      }
    
      private def readObject (in: java.io.ObjectInputStream): Unit = {
        conf = new Configuration()
        conf.readFields(in)
      }
    
      private def readObjectNoData(): Unit = {
        conf = new Configuration()
      }
    }
    

    请注意,对于所有 Writeable 类来说,将其设为通用是相对简单的;您只需要在构造函数中提供一个类名,并在反序列化期间使用它来实例化可写对象。

    【讨论】:

      【解决方案2】:

      您可以使用org.apache.spark.SerializableWritable 序列化和反序列化org.apache.hadoop.conf.Configuration

      例如:

      import org.apache.spark.SerializableWritable
      
      ...
      
      val hadoopConf = spark.sparkContext.hadoopConfiguration
      // serialize here
      val serializedConf = new SerializableWritable(hadoopConf)
      
      
      // then access the conf by calling .value on serializedConf
      rdd.map(someFunction(serializedConf.value))
      
      

      【讨论】:

        【解决方案3】:

        根据@Steve 的回答,这是一个 java 实现。

        import java.io.Serializable;
        import java.io.IOException;
        import org.apache.hadoop.conf.Configuration;
        
        
        public class SerializableHadoopConfiguration implements Serializable {
            Configuration conf;
        
            public SerializableHadoopConfiguration(Configuration hadoopConf) {
                this.conf = hadoopConf;
        
                if (this.conf == null) {
                    this.conf = new Configuration();
                }
            }
        
            public SerializableHadoopConfiguration() {
                this.conf = new Configuration();
            }
        
            public Configuration get() {
                return this.conf;
            }
        
            private void writeObject(java.io.ObjectOutputStream out) throws IOException {
                this.conf.write(out);
            }
        
            private void readObject(java.io.ObjectInputStream in) throws IOException {
                this.conf = new Configuration();
                this.conf.readFields(in);
            }
        }
        

        【讨论】:

          【解决方案4】:

          好像做不到,所以这里是我用的代码:

          final hdfsNameNodePath = "hdfs://quickstart.cloudera:8080";
          
          JavaPairRDD<String, PortableDataStream> imageByteRDD = jsc.binaryFiles(sourcePath);
                  if(!imageByteRDD.isEmpty())
                      imageByteRDD.foreachPartition(new VoidFunction<Iterator<Tuple2<String,PortableDataStream>>>() {
          
                          @Override
                          public void call(Iterator<Tuple2<String, PortableDataStream>> arg0)
                                  throws Exception {
          
                              Configuration conf = new Configuration();
                              conf.set("fs.defaultFS", hdfsNameNodePath);
                              //the string above should be passed as argument
          SequenceFile.Writer writer = SequenceFile.createWriter(
                                               conf, 
                                               SequenceFile.Writer.file([***ETCETERA...
          

          【讨论】:

            【解决方案5】:

            查看 Spark 内部代码库,应该广播 Hadoop 配置的序列化版本。

            https://github.com/apache/spark/blob/5d45a415f3a29898d92380380cfd82bfc7f579ea/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/binaryfile/BinaryFileFormat.scala#L98-L99

            val spark = SparkSession.builder.master("local").getOrCreate
            val broadcastedHadoopConf = spark.sparkContext.broadcast(new org.apache.spark.util.SerializableConfiguration(spark.sparkContext.hadoopConfiguration))
            
            val dfFiles = spark.read.format("binaryFile").load("/somepath").select("path")
            
            val df = dfFiles.map {row => {
              val rawPath = row.getString(0)
              val path = new Path(new URI(rawPath.replace(" ", "%20")))
            
              // get hadoop configuration in RDD method
              val hadoopConf = broadcastedHadoopConf.value.value
            
              val fs = path.getFileSystem(hadoopConfiguration)
              val status = fs.getFileStatus(path)
              val inputStream = fs.open(status.getPath)
              // ... whatever you need to do to read data
            }}
            

            【讨论】:

              猜你喜欢
              • 2018-08-12
              • 1970-01-01
              • 2016-08-21
              • 1970-01-01
              • 2018-01-04
              • 1970-01-01
              • 2020-07-25
              • 1970-01-01
              相关资源
              最近更新 更多