Spark newAPIHadoopFile读取Avro数据出现重复记录求助
Spark 1.6读取Avro记录时collect出现全量重复最后一条的问题解决
嘿,这个问题我之前也踩过坑,完全是Hadoop的RecordReader特性搞的鬼!咱们一步步理清楚:
问题根源
Hadoop的RecordReader为了减少对象创建的性能开销,会重用同一个Writable对象来处理所有记录。也就是说,每读取一条新记录,它只是把新数据写入到已有的对象实例里,而不是生成新的对象。
- 当你用
map(a -> a._1.datum())时,你只是把这个被重用的对象的引用放到了RDD中。等执行collect()的时候,所有元素其实都是指向同一个对象的引用,而这个对象最后被填充的是最后一条记录的数据——所以你看到全是重复的最后一条。 - 那为什么
foreach能正常打印?因为foreach是在每条记录刚被读取、对象还没被下一条数据覆盖的时候就执行了打印操作,这时候对象里存的是当前记录的内容。
解决方案:复制Avro记录
按照API文档的提示,必须在map阶段复制每个Avro记录,确保每个RDD元素都是独立的对象实例。
如果你的Avro是生成的具体类(Specific Record),可以用SpecificData的deepCopy方法:
import org.apache.avro.specific.SpecificData; JavaPairRDD<AvroKey, NullWritable> testRDD = sc.newAPIHadoopFile( dataSource.getFileDir(), AvroKeyInputFormat.class, AvroKey.class, NullWritable.class, new Configuration() ); List<Object> collect = testRDD.map(record -> { Object originalDatum = record._1.datum(); // 深度复制生成新的独立实例 return SpecificData.get().deepCopy(originalDatum.getSchema(), originalDatum); }).collect();
如果是通用记录(Generic Record),则用GenericData的deepCopy:
import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; List<Object> collect = testRDD.map(record -> { GenericRecord original = (GenericRecord) record._1.datum(); return GenericData.get().deepCopy(original.getSchema(), original); }).collect();
这样处理后,每个RDD元素都是独立的对象,collect到Driver端时就不会出现被覆盖的情况了。
内容的提问来源于stack exchange,提问作者karoma
相关产品推荐
相关产品推荐

