【问题标题】:reading/writing avro file in spark core using java使用java在spark core中读/写avro文件
【发布时间】:2016-08-12 06:02:10
【问题描述】:

我需要在 Spark 核心上用 java 编写的程序中访问 avro 文件数据。我可以使用 MapReduce InputFormat 类,但它给了我一个包含每一行文件作为键的元组。因为我没有使用 scala,所以很难解析它。

JavaPairRDD<AvroKey<GenericRecord>, AvroValue> avroRDD = sc.newAPIHadoopFile("dataset/testfile.avro", AvroKeyInputFormat.class, AvroKey.class, NullWritable.class,new Configuration()); 

是否有可用的实用程序类或 jar 可用于将 avro 数据直接映射到 java 类。例如。 codehaus.jackson 包提供了将 json 映射到 java 类的规定。

否则是否有任何其他方法可以轻松地将 avro 文件中存在的字段解析为 java 类或 RDD。

【问题讨论】:

    标签: java apache-spark avro


    【解决方案1】:

    考虑您的 avro 文件包含序列化对,键是 String,值是 avro 类。然后你可以拥有一些Utils 类的通用静态函数,如下所示:

    public class Utils {
    
      public static <T> JavaPairRDD<String, T> loadAvroFile(JavaSparkContext sc, String avroPath) {
        JavaPairRDD<AvroKey, NullWritable> records = sc.newAPIHadoopFile(avroPath, AvroKeyInputFormat.class, AvroKey.class, NullWritable.class, sc.hadoopConfiguration());
        return records.keys()
            .map(x -> (GenericRecord) x.datum())
            .mapToPair(pair -> new Tuple2<>((String) pair.get("key"), (T)pair.get("value")));
      }
    }
    

    然后你可以这样使用该方法:

    JavaPairRDD<String, YourAvroClassName> records = Utils.<YourAvroClassName>loadAvroFile(sc, inputDir);
    

    您可能还需要使用KryoSerializer 并注册您的自定义KryoRegistrator

    sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer");
    sparkConf.set("spark.kryo.registrator", "com.test.avro.MyKryoRegistrator");
    

    注册器类看起来是这样的:

    public class MyKryoRegistrator implements KryoRegistrator {
    
      public static class SpecificInstanceCollectionSerializer<T extends Collection> extends CollectionSerializer {
        Class<T> type;
        public SpecificInstanceCollectionSerializer(Class<T> type) {
          this.type = type;
        }
    
        @Override
        protected Collection create(Kryo kryo, Input input, Class<Collection> type) {
          return kryo.newInstance(this.type);
        }
    
        @Override
        protected Collection createCopy(Kryo kryo, Collection original) {
          return kryo.newInstance(this.type);
        }
      }
    
    
      Logger logger = LoggerFactory.getLogger(this.getClass());
    
      @Override
      public void registerClasses(Kryo kryo) {
        // Avro POJOs contain java.util.List which have GenericData.Array as their runtime type
        // because Kryo is not able to serialize them properly, we use this serializer for them
        kryo.register(GenericData.Array.class, new SpecificInstanceCollectionSerializer<>(ArrayList.class));
        kryo.register(YourAvroClassName.class);
      }
    }
    

    【讨论】:

      猜你喜欢
      • 2016-05-02
      • 2018-01-03
      • 1970-01-01
      • 2015-10-31
      • 2014-01-03
      • 1970-01-01
      • 2019-05-11
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多