【发布时间】: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 等价物更感兴趣。感谢您的评论!
标签: java apache-spark hdfs avro spark-avro