【问题标题】:How to serialize the data to AVRO schema in Spark (with Java)?如何在 Spark(使用 Java)中将数据序列化为 AVRO 模式?
【发布时间】:2016-08-01 12:20:21
【问题描述】:

我已经定义了一个 AVRO 模式,并为这些模式生成了一些带有 avro-tools 的类。现在,我想将数据序列化到磁盘。我为此找到了一些关于 scala 的答案,但不适用于 Java。 Article 类是使用 avro-tools 生成的,是由我定义的模式制成的。

这是我尝试如何做的代码的简化版本:

JavaPairRDD<String, String> filesRDD = context.wholeTextFiles(inputDataPath);
JavaRDD<Article> processingFiles = filesRDD.map(fileNameContent -> {
    // The name of the file
    String fileName = fileNameContent._1();
    // The content of the file
    String fileContent = fileNameContent._2();

    // An object from my avro schema
    Article a = new Article(fileContent);

    Processing processing = new Processing();
    // .... some processing of the content here ... //

    processing.serializeArticleToDisk(avroFileName);

    return a;
});

其中serializeArticleToDisk(avroFileName)定义如下:

public void serializeArticleToDisk(String filename) throws IOException{
    // Serialize article to disk
    DatumWriter<Article> articleDatumWriter = new SpecificDatumWriter<Article>(Article.class);
    DataFileWriter<Article> dataFileWriter = new DataFileWriter<Article>(articleDatumWriter);
    dataFileWriter.create(this.article.getSchema(), new File(filename));
    dataFileWriter.append(this.article);
    dataFileWriter.close();
}

Article 是我的 avro 架构。

现在,映射器向我抛出错误:

java.io.FileNotFoundException: hdfs:/...path.../avroFileName.avro (No such file or directory)   
at java.io.FileOutputStream.open0(Native Method)    
at java.io.FileOutputStream.open(FileOutputStream.java:270)     
at java.io.FileOutputStream.<init>(FileOutputStream.java:213)   
at java.io.FileOutputStream.<init>(FileOutputStream.java:162)   
at org.apache.avro.file.SyncableFileOutputStream.<init>(SyncableFileOutputStream.java:60)   
at org.apache.avro.file.DataFileWriter.create(DataFileWriter.java:129)
at org.apache.avro.file.DataFileWriter.create(DataFileWriter.java:129)
at sentences.ProcessXML.serializeArticleToDisk(ProcessXML.java:207)     
. . . rest of the stacktrace ... 

虽然文件路径是正确的。

之后我使用了collect() 方法,因此map 函数中的其他所有内容都可以正常工作(序列化部分除外)。

我对 Spark 很陌生,所以我不确定这实际上是否微不足道。我怀疑我需要使用一些写入功能,而不是在映射器中进行写入(但不确定这是否属实)。任何想法如何解决这个问题?

编辑:

我显示的错误堆栈跟踪的最后一行,实际上是在这部分:

dataFileWriter.create(this.article.getSchema(), new File(filename));

这是引发实际错误的部分。我假设dataFileWriter 需要替换为其他内容。有什么想法吗?

【问题讨论】:

标签: java apache-spark hdfs avro spark-avro


【解决方案1】:

看来你使用 Spark 的方式不对。

Map 是一个转换函数。仅调用map 不会调用RDD 的计算。您必须调用 action,例如 forEach()collect()

另请注意,提供给map 的 lambda 将在驱动程序中序列化并传输到集群中的某些Node

添加

尝试使用 Spark SQL 和 Spark-Avro 将 Spark DataFrame 保存为 Avro 格式:

// Load a text file and convert each line to a JavaBean.
JavaRDD<Person> people = sc.textFile("/examples/people.txt")
    .map(Person::parse);

// Apply a schema to an RDD
DataFrame peopleDF = sqlContext.createDataFrame(people, Person.class);
peopleDF.write()
    .format("com.databricks.spark.avro")
    .save("/output");

【讨论】:

  • 你在说什么——map 绝对会调用RDD 的计算。 map 返回一个新的RDD,其中所有元素都基于map 函数重新计算。
  • @Denis Kokorin:之后我正在使用collect(),所以map 中的所有内容都已经正常工作了,这很好。除了序列化之外的任何东西都可以在 map 函数中使用。
  • 也许他的意思是你应该在map 后面加上foreach 并在那里写信?如果这个答案有示例代码可能会有所帮助。
  • 抱歉之前没有提到。但是你的错误指向hdfs:/...path.../avroFileName.avro。默认情况下,Java 不解析 HDFS 协议。尝试使用 Hadoop 的 FileSystem 打开OutputStream。此外,您绝对不应该使用map() 将某些内容保存到 HDFS。使用foreach()store()
  • 我已经编辑了我的原始帖子。很抱歉误导了你。我匆匆写下这个答案。
【解决方案2】:

此解决方案不使用数据帧,也没有抛出任何错误:

import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.io.NullWritable;
import org.apache.avro.mapred.AvroKey;
import org.apache.spark.api.java.JavaPairRDD;
import scala.Tuple2;

   .  .  .  .  .

// Serializing to AVRO
JavaPairRDD<AvroKey<Article>, NullWritable> javaPairRDD = processingFiles.mapToPair(r -> {    
    return new Tuple2<AvroKey<Article>, NullWritable>(new AvroKey<Article>(r), NullWritable.get());
});
Job job = AvroUtils.getJobOutputKeyAvroSchema(Article.getClassSchema());
javaPairRDD.saveAsNewAPIHadoopFile(outputDataPath, AvroKey.class, NullWritable.class, AvroKeyOutputFormat.class, 
        job.getConfiguration());

AvroUtils.getJobOutputKeyAvroSchema 在哪里:

public static Job getJobOutputKeyAvroSchema(Schema avroSchema) {
    Job job;

    try {
        job = new Job();
    } catch (IOException e) {
        throw new RuntimeException(e);
    }

    AvroJob.setOutputKeySchema(job, avroSchema);
    return job;
}

Spark + Avro 的类似功能可以在这里找到 -> https://github.com/CeON/spark-utils

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-16
    • 1970-01-01
    • 2015-12-22
    • 2012-10-26
    • 2020-03-04
    • 2021-06-20
    相关资源
    最近更新 更多