You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:23:30