【问题标题】:coalesce does not reduce my number of output files合并不会减少我的输出文件数量
【发布时间】:2015-05-05 06:25:50
【问题描述】:

我有一个管理 RDD[SpecificRecordBase] 的 spark 工作,在 HDFS 上。

我的问题是它会生成很多文件,包括 95% 的空 avro 文件。 我尝试使用 coalesce 来减少我的 RDD 上的分区数量,从而减少我的输出文件的数量,但它没有效果。

 def write(data: RDD[SpecificRecordBase]) = {
   data.coalesce(1, false)    //has no effect
   val conf = new Configuration()
   val job = new org.apache.hadoop.mapreduce.Job(conf)

   AvroJob.setOutputKeySchema(job, schema)
   val pair = new PairRDDFunctions(rdd)
   pair.saveAsNewAPIHadoopFile(
     outputAvroDataPath,
     classOf[AvroKey[SpecificRecordBase]],
     classOf[org.apache.hadoop.io.NullWritable],
     classOf[AvroKeyOutputFormat[SpecificRecordBase]],
     job.getConfiguration)
}

我想rdd 分区配置和HDFS 分区之间丢失了一些东西,也许saveAsNewAPIHadoopFile 没有考虑到它,但我不确定。

我错过了什么吗?

有人能解释一下根据 rdd 分区调用saveAsNewAPIHadoopFile 时真正附加的内容吗?

【问题讨论】:

  • 您确定要输出data RDD?我没有从你的代码中看到它
  • 对不起,我忘记了一行代码:
  • 可以添加吗?
  • 是的,确实,我在这个示例中忘记了一行代码:val rdd = data.map(t => (new AvroKey(t), org.apache.hadoop.io.NullWritable.get)),是的,我应该在 rdd val 而不是数据上应用合并!然后,coalesce 返回一个新的 RDD,而不是更新当前的 RDD,所以它应该是:val rdd = data.map(t => (new AvroKey(t), org.apache.hadoop.io.NullWritable.get)).coalesce(1, false) 我的错误...谢谢你的代码审查 :)

标签: output apache-spark avro coalesce rdd


【解决方案1】:

感谢@0x0FFF 回答我自己的问题,正确的代码应该是:

    def write(data: RDD[SpecificRecordBase]) = {
           val rdd = data.map(t => (new AvroKey(t), org.apache.hadoop.io.NullWritable.get))
           val rdd1Partition = rdd.coalesce(1, false)  //change nb of partitions to 1

           val conf = new Configuration()
           val job = new org.apache.hadoop.mapreduce.Job(conf)

           AvroJob.setOutputKeySchema(job, schema)
           val pair = new PairRDDFunctions(rdd1Partition) //so only one file will be in output
           pair.saveAsNewAPIHadoopFile(
             outputAvroDataPath,
             classOf[AvroKey[SpecificRecordBase]],
             classOf[org.apache.hadoop.io.NullWritable],
             classOf[AvroKeyOutputFormat[SpecificRecordBase]],
             job.getConfiguration)
        }

再次感谢您!

【讨论】:

  • 澄清更改:Seb 使用了在使用 coalesce() 后创建的新 rdd(rdd1Partition)。就像继续使用旧的 rdd 变量(rdd)一样,它在执行时不会包括沿袭中的合并。
猜你喜欢
  • 2011-08-07
  • 1970-01-01
  • 1970-01-01
  • 2022-06-29
  • 2014-11-07
  • 1970-01-01
  • 1970-01-01
  • 2019-09-01
  • 1970-01-01
相关资源
最近更新 更多