【问题标题】:External command s3-dist-cp execution in spark-scala through scala.sys.process API通过 scala.sys.process API 在 spark-scala 中执行外部命令 s3-dist-cp
【发布时间】:2019-06-02 16:05:00
【问题描述】:

当我在 unix shell/终端中运行所有这 3 个命令时,它们都可以正常工作,将退出状态返回为 0

unix_shell> ls -la
unix_shell> hadoop fs -ls /user/hadoop/temp
unix_shell> s3-dist-cp --src ./abc.txt --dest s3://bucket/folder/

现在我正在尝试通过 scala 进程 api 作为外部进程运行这些相同的命令,示例代码如下:

import scala.sys.process._

val cmd_1 = "ls -la"
val cmd_2 = "hadoop fs -ls /user/hadoop/temp/"
val cmd_3 = "/usr/bin/s3-dist-cp --src /tmp/sample.txt --dest s3://bucket/folder/"
val cmd_4 = "s3-dist-cp --src /tmp/sample.txt --dest s3://bucket/folder/"

val exitCode_1 = (stringToProcess(cmd_1)).! // works fine and produces result
val exitCode_2 = (stringToProcess(cmd_2)).! // works fine and produces result
val exitCode_3 = (stringToProcess(cmd_3)).! // **it just hangs, yielding nothing**
val exitCode_4 = (stringToProcess(cmd_4)).! // **it just hangs, yielding nothing**

上面 cmd_3 和 cmd_4 的区别只是绝对路径。 我在 spark-submit 脚本中明确传递了相关依赖项,如下所示

--jars hdfs:///user/hadoop/s3-dist-cp.jar

您的意见/建议会有所帮助。谢谢!

【问题讨论】:

  • 您使用的是集群模式还是客户端模式?如果集群是你复制的工人?在这种情况下,如果可行,请先尝试客户端模式,然后尝试集群模式。
  • 我只在客户端模式下运行,我只打算在 spark 驱动节点本身上运行 s3-dist-cp 命令,
  • 关于所有 hadoop 应用程序 jar 的任何想法,我需要包含到 --jars 选项,我尝试使用 hdfs:///user/hadoop/s3-dist-cp.jar 但不能工作,这个 jar 是否足够,或者我需要在驱动程序类路径中包含更多 jar

标签: java scala amazon-web-services apache-spark amazon-s3


【解决方案1】:

看来你做的是对的。看这里 https://github.com/gorros/spark-scala-tips/blob/master/README.md

import scala.sys.process._

def s3distCp(src: String, dest: String): Unit = {
    s"s3-dist-cp --src $src --dest $dest".!
}

请检查此注释...我想知道您是否属于这种情况。

关于您的--jars /usr/lib/hadoop/client/*.jar

您可以使用tr 命令添加与s3-dist-cp 相关的jar,例如this. see my answer

--jars $(echo /dir_of_jars/*.jar | tr ' ' ',')

注意:为了能够使用此方法,您需要添加 Hadoop 应用程序,并且您需要在客户端或本地模式下运行 Spark,因为 s3-dist-cp 在从节点上不可用。如果要在集群模式下运行,则在引导期间将s3-dist-cp 命令复制到从属。

【讨论】:

  • 感谢 Ram 的回复,是的,我仅在客户端模式下使用 spark,并且我打算仅在客户端模式下的 spark 驱动程序节点本身上运行此命令。 “你需要添加 Hadoop 应用程序”——这是否意味着我需要包含 hadoop-aws.jar 或 hadoop-client。 jar,或者任何我需要精确包含的所有jar集的想法,但是我正在尝试这个解决方案,如果我有任何突破,会告诉你,
  • --jars /usr/lib/hadoop/client/*.jar,/usr/lib/hadoop/*.jar 我提供所有的 jars 并且这些 jars 出现在驱动程序类路径中然而 s3-dist-cp 挂起,不知道还需要做什么,当我直接从机器运行 spark-submit 命令时,它在完成地图作业 100% 后仍然挂起,似乎没有开始减少
【解决方案2】:

实际上 scala 进程在 spark 上下文之外工作,所以为了成功运行 s3-dist-cp 命令,我所做的只是在启动包装 s3-dist-cp 命令的 scala 进程之前停止 spark 上下文,完整的工作代码如下:

    logger.info("Moving ORC files from HDFS to S3 !!")

    import scala.sys.process._

    logger.info("stopping spark context..##")
    val spark = IngestionContext.sparkSession
    spark.stop()
    logger.info("spark context stopped..##")
    logger.info("sleeping for 10 secs")
    Thread.sleep(10000) // this sleep is not required, this was just for debugging purpose, you can remove this in your final code.
    logger.info("woke up after sleeping for 10 secs")

    try {
      /**
       * following is the java version, off course you need take care of few imports
       */
      //val pb = new java.lang.ProcessBuilder("s3-dist-cp", "--src", INGESTED_ORC_DIR, "--dest", "s3:/" + paramMap(Storage_Output_Path).substring(4) + "_temp", "--srcPattern", ".*\\.orc")
      //val pb = new java.lang.ProcessBuilder("hadoop", "jar", "/usr/share/aws/emr/s3-dist-cp/lib/s3-dist-cp.jar", "--src", INGESTED_ORC_DIR, "--dest", "s3:/" + paramMap(Storage_Output_Path).substring(4) + "_temp", "--srcPattern", ".*\\.orc")
      //pb.directory(new File("/tmp"))
      //pb.inheritIO()
      //pb.redirectErrorStream(true)
      //val process = pb.start()
      //val is = process.getInputStream()
      //val isr = new InputStreamReader(is)
      //val br = new BufferedReader(isr)
      //var line = ""
      //logger.info("printling lines:")
      //while (line != null) {
      //  line = br.readLine()
      //  logger.info("line=[{}]", line)
      //}

      //logger.info("process goes into waiting state")
      //logger.info("Waited for: " + process.waitFor())
      //logger.info("Program terminated!")

      /**
       * following is the scala version
       */
      val S3_DIST_CP = "s3-dist-cp"
      val INGESTED_ORC_DIR = S3Util.getSaveOrcPath()

      // listing out all the files
      //val s3DistCpCmd = S3_DIST_CP + " --src " + INGESTED_ORC_DIR + " --dest " + paramMap(Storage_Output_Path).substring(4) + "_temp --srcPattern .*\\.orc"
      //-Dmapred.child.java.opts=-Xmx1024m -Dmapreduce.job.reduces=2
      val cmd = S3_DIST_CP + " --src " + INGESTED_ORC_DIR + " --dest " + "s3:/" + paramMap(Storage_Output_Path).substring(4) + "_temp --srcPattern .*\\.orc"

      //val cmd = "hdfs dfs -cp -f " + INGESTED_ORC_DIR + "/* " + "s3:/" + paramMap(Storage_Output_Path).substring(4) + "_temp/"
      //val cmd = "hadoop distcp " + INGESTED_ORC_DIR + "/ s3:/" + paramMap(Storage_Output_Path).substring(4) + "_temp_2/"

      logger.info("full hdfs to s3 command : [{}]", cmd)

      // command execution
      val exitCode = (stringToProcess(cmd)).!

      logger.info("s3_dist_cp command exit code: {} and s3 copy got " + (if (exitCode == 0) "SUCCEEDED" else "FAILED"), exitCode)
    } catch {
      case ex: Exception =>
        logger.error(
          "there was an exception while copying orc file to s3 bucket. {} {}",
          "", ex.getMessage, ex)
        throw new IngestionException("s3 dist cp command failure", null, Some(StatusEnum.S3_DIST_CP_CMD_FAILED))
    }

虽然上面的代码完全按照预期工作,但也有其他观察结果,如下所示:

而不是使用这个

val exitCode = (stringToProcess(cmd)).!

如果你使用这个

val exitCode = (stringToProcess(cmd)).!!

注意 single 的区别!和双!!,作为单!只返回退出代码,而 double !!返回流程执行的输出

所以在单身的情况下!上面的代码运行良好,如果是 double !!,它也可以运行,但是它在 S3 存储桶中生成了太多的文件和副本,而不是原始文件的数量。

至于 spark-submit 命令,无需担心 --driver-class-path 甚至 --jars 选项,因为我没有传递任何依赖项。

【讨论】:

  • 非常感谢@Ram Ghadiyaram,虽然我尝试了您的解决方案但没有奏效,但如果您没有通过分享该链接给我一点点推动,我可能没有进一步努力并到达一个解决方案,非常感谢,继续帮助:)
  • 好兄弟! :-) +1 ...一件事是你没有分享代码来指出究竟出了什么问题。我什至不知道你没有停止火花上下文。下次分享代码 sn-p 和软件版本,以获得优雅/精英的答案。一切顺利!但总的来说,我的回答仍然有效。我为你做了很多研究并在那里发布了答案。
  • 哦,感谢您的努力。这真的很有帮助,是的,我会负责分享最佳代码片段,以便在未来有更多的洞察力,实际上我对 SO 还是新手并学习这些东西。谢谢!!
猜你喜欢
  • 1970-01-01
  • 2020-07-23
  • 2018-09-16
  • 1970-01-01
  • 2019-08-28
  • 1970-01-01
  • 1970-01-01
  • 2013-04-16
  • 1970-01-01
相关资源
最近更新 更多