一次,我需要将一个hive表进行简单处理后写入另一个hive表,hive表数据我用的存储介质是腾讯云的cos对象存储。数据数据量很大其实也还好,不是很夸张,就是分区较多几千个吧,然后上游任务也没有控制好小文件,导致每个分区里都要好多小文件,在此背景下,任务执行了好几个小时了,我看driver和executor都还在运行没有停下来意思,那不用多想,肯定是哪儿出了问题,按照我排除数据任务问题的思路,先去看看driver的日志,没有啥异常,日志老早就不打了,再看看executor的日志,一样,日志老早也不打了,并且executor的日志最后一行显示的是最后的task执行成功并反馈到了driver。
在这里插入图片描述
webUI上也是执行完成了
在这里插入图片描述
这里没发现问题,然后上driver pod里看看线程快照,没看出个所以然来,就看到主线程在running,但是read操作,以为它是其它读数据的连接没释放而已,还在苦苦寻找是否有线程死锁等问题,后来在腾讯云大佬的帮助下(这就是用云产品的一个好处吧),这里是数据确实是已经写完落盘了,但是是落盘到临时目录,我说这儿我知道,写完到临时目录成功后,把临时目录的数据移出来就ok了,我还在认为是hive服务那边建表有问题导致任务一致完成不了,然后大佬说主线程这儿正是在做临时目录数据移动到正式文件目录里,我说为啥这么慢啊,大佬解释道,文件移动是一个rename的操作,本身在hdfs上是很快的,但因为这是对象存储,对象存储没有直接rename的操作,只能 list+copy+delete 的方式去实现,那么我这么多executor执行的task都花了这么长时间,那么rename在driver端单线程处理的话,花费的时间更是无法预估的长。。。。。然后大佬建议我将spark写hdfs的filecommitter的版本改为2,我这才了解到 filecommitter。
在这里插入图片描述
在这里插入图片描述
那么,接下来我们来了解一下这个spark写hdfs的filecommiter吧。

public class FileOutputCommitter extends PathOutputCommitter {
    private static final Logger LOG = LoggerFactory.getLogger(FileOutputCommitter.class);
    public static final String PENDING_DIR_NAME = "_temporary";
    /** @deprecated */
    @Deprecated
    protected static final String TEMP_DIR_NAME = "_temporary";
    public static final String SUCCEEDED_FILE_NAME = "_SUCCESS";
    public static final String SUCCESSFUL_JOB_OUTPUT_DIR_MARKER = "mapreduce.fileoutputcommitter.marksuccessfuljobs";
    public static final String FILEOUTPUTCOMMITTER_ALGORITHM_VERSION = "mapreduce.fileoutputcommitter.algorithm.version";
    public static final int FILEOUTPUTCOMMITTER_ALGORITHM_VERSION_DEFAULT = 2;
    public static final String FILEOUTPUTCOMMITTER_CLEANUP_SKIPPED = "mapreduce.fileoutputcommitter.cleanup.skipped";
    public static final boolean FILEOUTPUTCOMMITTER_CLEANUP_SKIPPED_DEFAULT = false;
    public static final String FILEOUTPUTCOMMITTER_CLEANUP_FAILURES_IGNORED = "mapreduce.fileoutputcommitter.cleanup-failures.ignored";
    public static final boolean FILEOUTPUTCOMMITTER_CLEANUP_FAILURES_IGNORED_DEFAULT = false;
    public static final String FILEOUTPUTCOMMITTER_FAILURE_ATTEMPTS = "mapreduce.fileoutputcommitter.failures.attempts";
    public static final int FILEOUTPUTCOMMITTER_FAILURE_ATTEMPTS_DEFAULT = 1;
    private Path outputPath;
    private Path workPath;
    private final int algorithmVersion;
    private final boolean skipCleanup;
    private final boolean ignoreCleanupFailures;

    @Private
    public FileOutputCommitter(Path outputPath, JobContext context) throws IOException {
        super(outputPath, context);
        this.outputPath = null;
        this.workPath = null;
        Configuration conf = context.getConfiguration();
        this.algorithmVersion = conf.getInt("mapreduce.fileoutputcommitter.algorithm.version", 2);
        LOG.info("File Output Committer Algorithm version is " + this.algorithmVersion);
        if (this.algorithmVersion != 1 && this.algorithmVersion != 2) {
            throw new IOException("Only 1 or 2 algorithm version is supported");
        } else {
            this.skipCleanup = conf.getBoolean("mapreduce.fileoutputcommitter.cleanup.skipped", false);
            this.ignoreCleanupFailures = conf.getBoolean("mapreduce.fileoutputcommitter.cleanup-failures.ignored", false);
            LOG.info("FileOutputCommitter skip cleanup _temporary folders under output directory:" + this.skipCleanup + ", ignore cleanup failures: " + this.ignoreCleanupFailures);
            if (outputPath != null) {
                FileSystem fs = outputPath.getFileSystem(context.getConfiguration());
                this.outputPath = fs.makeQualified(outputPath);
            }

        }
    }

    
    private void mergePaths(FileSystem fs, FileStatus from, Path to) throws IOException {
        if (LOG.isDebugEnabled()) {
            LOG.debug("Merging data from " + from + " to " + to);
        }

        FileStatus toStat;
        try {
            toStat = fs.getFileStatus(to);
        } catch (FileNotFoundException var10) {
            toStat = null;
        }

        if (from.isFile()) {
            if (toStat != null && !fs.delete(to, true)) {
                throw new IOException("Failed to delete " + to);
            }

            if (!fs.rename(from.getPath(), to)) {
                throw new IOException("Failed to rename " + from + " to " + to);
            }
        } else if (from.isDirectory()) {
            if (toStat != null) {
                if (!toStat.isDirectory()) {
                    if (!fs.delete(to, true)) {
                        throw new IOException("Failed to delete " + to);
                    }

                    this.renameOrMerge(fs, from, to);
                } else {
                    FileStatus[] var5 = fs.listStatus(from.getPath());
                    int var6 = var5.length;

                    for(int var7 = 0; var7 < var6; ++var7) {
                        FileStatus subFrom = var5[var7];
                        Path subTo = new Path(to, subFrom.getPath().getName());
                        this.mergePaths(fs, subFrom, subTo);
                    }
                }
            } else {
                this.renameOrMerge(fs, from, to);
            }
        }

    }

    private void renameOrMerge(FileSystem fs, FileStatus from, Path to) throws IOException {
        if (this.algorithmVersion == 1) {
            if (!fs.rename(from.getPath(), to)) {
                throw new IOException("Failed to rename " + from + " to " + to);
            }
        } else {
            fs.mkdirs(to);
            FileStatus[] var4 = fs.listStatus(from.getPath());
            int var5 = var4.length;

            for(int var6 = 0; var6 < var5; ++var6) {
                FileStatus subFrom = var4[var6];
                Path subTo = new Path(to, subFrom.getPath().getName());
                this.mergePaths(fs, subFrom, subTo);
            }
        }

    }

      @Private
  public void commitTask(TaskAttemptContext context, Path taskAttemptPath) 
      throws IOException {

    TaskAttemptID attemptId = context.getTaskAttemptID();
    if (hasOutputPath()) {
      context.progress();
      if(taskAttemptPath == null) {
        taskAttemptPath = getTaskAttemptPath(context);
      }
      FileSystem fs = taskAttemptPath.getFileSystem(context.getConfiguration());
      FileStatus taskAttemptDirStatus;
      try {
        taskAttemptDirStatus = fs.getFileStatus(taskAttemptPath);
      } catch (FileNotFoundException e) {
        taskAttemptDirStatus = null;
      }

      if (taskAttemptDirStatus != null) {
        if (algorithmVersion == 1) {
          Path committedTaskPath = getCommittedTaskPath(context);
          if (fs.exists(committedTaskPath)) {
             if (!fs.delete(committedTaskPath, true)) {
               throw new IOException("Could not delete " + committedTaskPath);
             }
          }
          if (!fs.rename(taskAttemptPath, committedTaskPath)) {
            throw new IOException("Could not rename " + taskAttemptPath + " to "
                + committedTaskPath);
          }
          LOG.info("Saved output of task '" + attemptId + "' to " +
              committedTaskPath);
        } else {
          // directly merge everything from taskAttemptPath to output directory
          mergePaths(fs, taskAttemptDirStatus, outputPath);
          LOG.info("Saved output of task '" + attemptId + "' to " +
              outputPath);
        }
      } else {
        LOG.warn("No Output found for " + attemptId);
      }
    } else {
      LOG.warn("Output Path is null in commitTask()");
    }
  }

  public void commitJob(JobContext context) throws IOException {
    int maxAttemptsOnFailure = isCommitJobRepeatable(context) ?
        context.getConfiguration().getInt(FILEOUTPUTCOMMITTER_FAILURE_ATTEMPTS,
            FILEOUTPUTCOMMITTER_FAILURE_ATTEMPTS_DEFAULT) : 1;
    int attempt = 0;
    boolean jobCommitNotFinished = true;
    while (jobCommitNotFinished) {
      try {
        commitJobInternal(context);
        jobCommitNotFinished = false;
      } catch (Exception e) {
        if (++attempt >= maxAttemptsOnFailure) {
          throw e;
        } else {
          LOG.warn("Exception get thrown in job commit, retry (" + attempt +
              ") time.", e);
        }
      }
    }
  }

    
 /**
   * The job has completed, so do following commit job, include:
   * Move all committed tasks to the final output dir (algorithm 1 only).
   * Delete the temporary directory, including all of the work directories.
   * Create a _SUCCESS file to make it as successful.
   * @param context the job's context
   */
  @VisibleForTesting
  protected void commitJobInternal(JobContext context) throws IOException {
    if (hasOutputPath()) {
      Path finalOutput = getOutputPath();
      FileSystem fs = finalOutput.getFileSystem(context.getConfiguration());

      if (algorithmVersion == 1) {
        for (FileStatus stat: getAllCommittedTaskPaths(context)) {
          mergePaths(fs, stat, finalOutput);
        }
      }

      if (skipCleanup) {
        LOG.info("Skip cleanup the _temporary folders under job's output " +
            "directory in commitJob.");
      } else {
        // delete the _temporary folder and create a _done file in the o/p
        // folder
        try {
          cleanupJob(context);
        } catch (IOException e) {
          if (ignoreCleanupFailures) {
            // swallow exceptions in cleanup as user configure to make sure
            // commitJob could be success even when cleanup get failure.
            LOG.error("Error in cleanup job, manually cleanup is needed.", e);
          } else {
            // throw back exception to fail commitJob.
            throw e;
          }
        }
      }
      // True if the job requires output.dir marked on successful job.
      // Note that by default it is set to true.
      if (context.getConfiguration().getBoolean(
          SUCCESSFUL_JOB_OUTPUT_DIR_MARKER, true)) {
        Path markerPath = new Path(outputPath, SUCCEEDED_FILE_NAME);
        // If job commit is repeatable and previous/another AM could write
        // mark file already, we need to set overwritten to be true explicitly
        // in case other FS implementations don't overwritten by default.
        if (isCommitJobRepeatable(context)) {
          fs.create(markerPath, true).close();
        } else {
          fs.create(markerPath).close();
        }
      }
    } else {
      LOG.warn("Output Path is null in commitJob()");
    }
  }

      @Override
  public boolean isCommitJobRepeatable(JobContext context) throws IOException {
    return algorithmVersion == 2;
  }

    
}

从上面这个源码(hadoop 3.0.0)可以看到,对应的参数名称为:mapreduce.fileoutputcommitter.algorithm.version;在对应的sparkConf里的key的话为:spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version, value为1或者2;
在上面源码的 commitTask方法中 ,version为1 的话,数据是从taskAttemptPath移到committedTaskPath的, version为2 的话 数据是从taskAttemptPath移到outputPath的;
在上面源码的commitJob方法中,不难发现,它是调用的 commitJobInternal方法,在commitJobInternal方法中,如果version为1 会额外进行一次 mergePaths 操作;
总结一下就是version=1 的话,commitTask是将临时目录的数据先移动到上一层的一个committedTaskPath目录,但不是最终的输出路径,但version=2的时候,commitTask时就将临时目录的数据直接移动到最终的输出路径;所有commitTask执行完成后,再由driver端单线程的进行commitJob,commitJob时,如果version=1,那么还得再进行一次数据的移动到最终目录,然后还有一些清理临时数据操作,标记成功等操作;这个时候就会突出两个的差别,version=1会进行两次移动数据的操作,而version=2只会进行一次,假如这个时候失败了,那么version=1还能继续重试,version=2 就不再进行重试了,因为它的maxAttemptsOnFailure=1,所以可能出现任务失败了,version=1 看不到数据,version=2能看到数据,但是这个数据是否正确保证不了。
最后,version=1,性能差一点,可靠性高;version=2,性能好一点,可靠性差一点;
回到最开始对象存储的rename,因为假如是基于磁盘存储的hdfs,进行rename或者merge 的开销还是可以接受的;但是对象存储的rename是的查找列表,复制,删除 三连的话,那么我选择version=2;虽然腾讯云提供了这些参数:但是用上感觉吧不太行呢,就没再深究这些参数的原理了。

    spark.hadoop.fs.cosn.committer.name: directory
    spark.hadoop.fs.cosn.committer.staging.conflict-mode: append
    spark.hadoop.fs.cosn.committer.threads: '8'
    spark.hadoop.mapreduce.outputcommitter.factory.scheme.cosn: com.qcloud.emr.fs.COSNCommitterFactoryAdaptor
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐