【问题标题】:Hadoop multiple outputs with speculative execution具有推测执行的 Hadoop 多输出
【发布时间】:2015-07-15 17:06:56
【问题描述】:

我有一个任务,它将 avro 输出写入由输入记录的几个字段组织的多个目录中。

例如 : 各国历年的流程记录 并写入国家/年的目录结构 例如: 输出/usa/2015/outputs_usa_2015.avro 输出/uk/2014/outputs_uk_2014.avro
AvroMultipleOutputs multipleOutputs=new AvroMultipleOutputs(context);
....
....
     multipleOutputs.write("output", avroKey, NullWritable.get(), 
            OUTPUT_DIR + "/" + record.getCountry() + "/" + record.getYear() + "/outputs_" +record.getCountry()+"_"+ record.getYear());

以下代码将使用哪个输出提交者来编写输出。与推测执行一起使用是否不安全? 通过推测执行,这会导致(可能导致)org.apache.hadoop.hdfs.server.namenode.LeaseExpiredException

在这篇文章中 Hadoop Reducer: How can I output to multiple directories using speculative execution? 建议使用自定义输出提交器

hadoop AvroMultipleOutputs 的以下代码没有说明推测执行有任何问题

 private synchronized RecordWriter getRecordWriter(TaskAttemptContext taskContext,
          String baseFileName) throws IOException, InterruptedException {

    writer =
                ((OutputFormat) ReflectionUtils.newInstance(taskContext.getOutputFormatClass(),
                    taskContext.getConfiguration())).getRecordWriter(taskContext);
...
}

如果 baseoutput 路径在作业目录之外,write 方法也不会记录任何问题

public void write(String namedOutput, Object key, Object value, String baseOutputPath)

在作业目录之外写入时,AvroMultipleOutputs(其他输出)是否存在推测执行的真正问题? 如果,那么我如何覆盖 AvroMultipleOutputs 以拥有它自己的输出提交者。我在 AvroMultipleOutputs 中看不到它使用其输出提交者的任何输出格式

【问题讨论】:

  • 你自己写实现了吗?我也有同样的问题。
  • 当您说“通过推测执行这会导致(可能导致)org.apache.hadoop.hdfs.server.namenode.LeaseExpiredException”时,您是否在任何地方看到过此文档,或者您是根据经验说话。我们看到了相同的行为,但没有找到任何明确的引用来禁用使用多个输出时的推测执行。
  • 是的,它已记录在案。这里有一个警告archive.cloudera.com/cdh5/cdh/5/hadoop/api/org/apache/hadoop/…

标签: java hadoop hadoop-yarn multipleoutputs speculative-execution


【解决方案1】:

AvroMultipleOutputs 将使用您在作业配置中注册的OutputFormat,同时添加命名输出,例如使用来自AvroMultipleOutputsaddNamedOutput API(例如AvroKeyValueOutputFormat)。

使用AvroMultipleOutputs,您可能无法使用推测任务执行功能。即使覆盖它也无济于事,也不会简单。

相反,您应该编写自己的OutputFormat(很可能扩展一种可用的 Avro 输出格式,例如AvroKeyValueOutputFormat),并覆盖/实现其getRecordWriter API,它会返回一个RecordWriter 实例说@ 987654332@(仅供参考)。

这个MainRecordWriter将维护RecordWriter(例如AvroKeyValueRecordWriter)实例的映射。这些RecordWriter 实例中的每一个都属于输出文件之一。在MainRecordWriterwrite API 中,您将从映射中获取实际的RecordWriter 实例(基于您要写入的记录),并使用此记录写入器写入记录。所以MainRecordWriter 只是作为多个RecordWriter 实例的包装器。

对于一些类似的实现,你可能想研究piggybank库中MultiStorage类的代码。

【讨论】:

    【解决方案2】:

    当您将命名输出添加到AvroMultipleOutputs 时,它将调用AvroKeyOutputFormat.getRecordWriter()AvroKeyValueOutputFormat.getRecordWriter(),后者调用AvroOutputFormatBase.getAvroFileOutputStream(),其内容为

    protected OutputStream getAvroFileOutputStream(TaskAttemptContext context) throws IOException {
      Path path = new Path(((FileOutputCommitter)getOutputCommitter(context)).getWorkPath(),
        getUniqueFile(context,context.getConfiguration().get("avro.mo.config.namedOutput","part"),org.apache.avro.mapred.AvroOutputFormat.EXT));
      return path.getFileSystem(context.getConfiguration()).create(path);
    }
    

    AvroOutputFormatBase 扩展FileOutputFormat(上述方法中的getOutputCommitter() 实际上是对FileOutputFormat.getOutputCommitter() 的调用。因此,AvroMultipleOutputs 应该具有与MultipleOutputs 相同的约束。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-07-19
      • 2013-01-27
      • 2013-02-16
      • 2014-06-09
      • 2013-10-01
      • 2019-01-16
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多