|
@@ -1,50 +1,123 @@
|
|
|
package com.cloudera.hue.sparker.yarn
|
|
package com.cloudera.hue.sparker.yarn
|
|
|
|
|
|
|
|
-import java.util.Collections
|
|
|
|
|
-
|
|
|
|
|
-import org.apache.hadoop.fs.Path
|
|
|
|
|
|
|
+import org.apache.hadoop.fs.{FileSystem, Path}
|
|
|
|
|
+import org.apache.hadoop.yarn.api.ApplicationConstants
|
|
|
|
|
+import org.apache.hadoop.yarn.api.ApplicationConstants.Environment
|
|
|
import org.apache.hadoop.yarn.api.records._
|
|
import org.apache.hadoop.yarn.api.records._
|
|
|
import org.apache.hadoop.yarn.client.api.YarnClient
|
|
import org.apache.hadoop.yarn.client.api.YarnClient
|
|
|
import org.apache.hadoop.yarn.conf.YarnConfiguration
|
|
import org.apache.hadoop.yarn.conf.YarnConfiguration
|
|
|
import org.apache.hadoop.yarn.util.{ConverterUtils, Records}
|
|
import org.apache.hadoop.yarn.util.{ConverterUtils, Records}
|
|
|
|
|
+import org.slf4j.LoggerFactory
|
|
|
|
|
|
|
|
import scala.collection.JavaConversions._
|
|
import scala.collection.JavaConversions._
|
|
|
|
|
|
|
|
-class Client {
|
|
|
|
|
- def main(args: Array[String]) = {
|
|
|
|
|
|
|
+object Client extends Logging {
|
|
|
|
|
+
|
|
|
|
|
+ def main(args: Array[String]): Unit = {
|
|
|
|
|
+ println(args.length)
|
|
|
|
|
+ args.foreach(println(_))
|
|
|
|
|
|
|
|
val packagePath = new Path(args(1))
|
|
val packagePath = new Path(args(1))
|
|
|
|
|
|
|
|
val yarnConf = new YarnConfiguration()
|
|
val yarnConf = new YarnConfiguration()
|
|
|
- val yarnClient = YarnClient.createYarnClient()
|
|
|
|
|
- yarnClient.init(yarnConf)
|
|
|
|
|
- yarnClient.start
|
|
|
|
|
|
|
+ val client = new Client(yarnConf)
|
|
|
|
|
|
|
|
try {
|
|
try {
|
|
|
- submitApplication(yarnClient, yarnConf, packagePath, List("echo hi"))
|
|
|
|
|
|
|
+ val job = client.submitApplication(
|
|
|
|
|
+ packagePath,
|
|
|
|
|
+ List(
|
|
|
|
|
+ "__package/bin/run-am.sh 1>%s/stdout 2>%s/stderr" format (
|
|
|
|
|
+ ApplicationConstants.LOG_DIR_EXPANSION_VAR,
|
|
|
|
|
+ ApplicationConstants.LOG_DIR_EXPANSION_VAR
|
|
|
|
|
+ )
|
|
|
|
|
+ /*
|
|
|
|
|
+ "/bin/pwd " +
|
|
|
|
|
+ " 1>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout" +
|
|
|
|
|
+ " 2>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr;",
|
|
|
|
|
+ "/bin/ls " +
|
|
|
|
|
+ " 1>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout" +
|
|
|
|
|
+ " 2>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr;",
|
|
|
|
|
+ "/bin/echo hi" +
|
|
|
|
|
+ " 1>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout" +
|
|
|
|
|
+ " 2>>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr"
|
|
|
|
|
+ */
|
|
|
|
|
+ )
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
|
|
+ info("waiting for job to start")
|
|
|
|
|
+
|
|
|
|
|
+ job.waitForStatus(Running(), 500) match {
|
|
|
|
|
+ case Some(Running()) => {
|
|
|
|
|
+ info("job started successfully")
|
|
|
|
|
+ }
|
|
|
|
|
+ case Some(appStatus) => {
|
|
|
|
|
+ warn("unable to start job successfully. job has status %s" format appStatus)
|
|
|
|
|
+ }
|
|
|
|
|
+ case None => {
|
|
|
|
|
+ warn("timed out waiting for job to start")
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ job.waitForFinish(100000) match {
|
|
|
|
|
+ case Some(SuccessfulFinish()) => {
|
|
|
|
|
+ info("job finished successfully")
|
|
|
|
|
+ }
|
|
|
|
|
+ case Some(appStatus) => {
|
|
|
|
|
+ info("job finished unsuccessfully %s" format appStatus)
|
|
|
|
|
+ }
|
|
|
|
|
+ case None => {
|
|
|
|
|
+ info("timed out")
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
} finally {
|
|
} finally {
|
|
|
- yarnClient.close()
|
|
|
|
|
|
|
+ client.close()
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+class Client(yarnConf: YarnConfiguration) {
|
|
|
|
|
+
|
|
|
|
|
+ import Client._
|
|
|
|
|
|
|
|
- def submitApplication(yarnClient: YarnClient, yarnConf: YarnConfiguration, packagePath: Path, cmds: List[String]) = {
|
|
|
|
|
|
|
+ val yarnClient = YarnClient.createYarnClient()
|
|
|
|
|
+ yarnClient.init(yarnConf)
|
|
|
|
|
+ yarnClient.start()
|
|
|
|
|
+
|
|
|
|
|
+ def submitApplication(packagePath: Path, cmds: List[String]): Job = {
|
|
|
val app = yarnClient.createApplication()
|
|
val app = yarnClient.createApplication()
|
|
|
val newAppResponse = app.getNewApplicationResponse
|
|
val newAppResponse = app.getNewApplicationResponse
|
|
|
|
|
|
|
|
val appId = newAppResponse.getApplicationId
|
|
val appId = newAppResponse.getApplicationId
|
|
|
|
|
|
|
|
|
|
+ info("preparing to submit %s" format appId)
|
|
|
|
|
+
|
|
|
val appContext = app.getApplicationSubmissionContext
|
|
val appContext = app.getApplicationSubmissionContext
|
|
|
|
|
+ appContext.setApplicationName(appId.toString)
|
|
|
|
|
+
|
|
|
val containerCtx = Records.newRecord(classOf[ContainerLaunchContext])
|
|
val containerCtx = Records.newRecord(classOf[ContainerLaunchContext])
|
|
|
val resource = Records.newRecord(classOf[Resource])
|
|
val resource = Records.newRecord(classOf[Resource])
|
|
|
- val packageResource = Records.newRecord(classOf[LocalResource])
|
|
|
|
|
|
|
|
|
|
- appContext.setApplicationName(appId.toString)
|
|
|
|
|
|
|
+ info("Copy app master jar from local filesystem and add to the local environment")
|
|
|
|
|
+ /*
|
|
|
|
|
+ val localResources = Map(
|
|
|
|
|
+ "app" => uploadLocalResource()
|
|
|
|
|
+ )
|
|
|
|
|
+ Map
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+ addToLocalResources(fs, appMasterJar, appMasterJarPath, appId, localResources, null)
|
|
|
|
|
+ */
|
|
|
|
|
+
|
|
|
|
|
+ val packageResource = Records.newRecord(classOf[LocalResource])
|
|
|
|
|
|
|
|
val packageUrl = ConverterUtils.getYarnUrlFromPath(packagePath)
|
|
val packageUrl = ConverterUtils.getYarnUrlFromPath(packagePath)
|
|
|
val fileStatus = packagePath.getFileSystem(yarnConf).getFileStatus(packagePath)
|
|
val fileStatus = packagePath.getFileSystem(yarnConf).getFileStatus(packagePath)
|
|
|
|
|
|
|
|
packageResource.setResource(packageUrl)
|
|
packageResource.setResource(packageUrl)
|
|
|
|
|
+ info("set package url to %s for %s" format (packageUrl, appId))
|
|
|
packageResource.setSize(fileStatus.getLen)
|
|
packageResource.setSize(fileStatus.getLen)
|
|
|
|
|
+ info("set package size to %s for %s" format (fileStatus.getLen, appId))
|
|
|
packageResource.setTimestamp(fileStatus.getModificationTime)
|
|
packageResource.setTimestamp(fileStatus.getModificationTime)
|
|
|
packageResource.setType(LocalResourceType.ARCHIVE)
|
|
packageResource.setType(LocalResourceType.ARCHIVE)
|
|
|
packageResource.setVisibility(LocalResourceVisibility.APPLICATION)
|
|
packageResource.setVisibility(LocalResourceVisibility.APPLICATION)
|
|
@@ -54,14 +127,99 @@ class Client {
|
|
|
appContext.setResource(resource)
|
|
appContext.setResource(resource)
|
|
|
|
|
|
|
|
containerCtx.setCommands(cmds)
|
|
containerCtx.setCommands(cmds)
|
|
|
- containerCtx.setLocalResources(Collections.singletonMap("__package", packageResource))
|
|
|
|
|
|
|
+ containerCtx.setLocalResources(Map("__package" -> packageResource))
|
|
|
|
|
+
|
|
|
appContext.setApplicationId(appId)
|
|
appContext.setApplicationId(appId)
|
|
|
appContext.setAMContainerSpec(containerCtx)
|
|
appContext.setAMContainerSpec(containerCtx)
|
|
|
appContext.setApplicationType("sparker")
|
|
appContext.setApplicationType("sparker")
|
|
|
|
|
+
|
|
|
|
|
+ info("submitting application request for %s" format appId)
|
|
|
|
|
+
|
|
|
yarnClient.submitApplication(appContext)
|
|
yarnClient.submitApplication(appContext)
|
|
|
|
|
+
|
|
|
|
|
+ new Job(yarnClient, appId)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def close(): Unit = {
|
|
|
|
|
+ yarnClient.close()
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private def addToLocalResources(fs: FileSystem, fileSrcPath: String, fileDstPath: String, appId: String): LocalResource = {
|
|
|
|
|
+ val appName = "sparker"
|
|
|
|
|
+ val suffix = appName + "/" + appId + "/" + fileDstPath
|
|
|
|
|
+
|
|
|
|
|
+ val dst = new Path(fs.getHomeDirectory, suffix)
|
|
|
|
|
+
|
|
|
|
|
+ fs.copyFromLocalFile(new Path(fileSrcPath), dst)
|
|
|
|
|
+
|
|
|
|
|
+ val srcFileStatus = fs.getFileStatus(dst)
|
|
|
|
|
+ LocalResource.newInstance(
|
|
|
|
|
+ ConverterUtils.getYarnUrlFromURI(dst.toUri),
|
|
|
|
|
+ LocalResourceType.FILE,
|
|
|
|
|
+ LocalResourceVisibility.APPLICATION,
|
|
|
|
|
+ srcFileStatus.getLen,
|
|
|
|
|
+ srcFileStatus.getModificationTime
|
|
|
|
|
+ )
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+class Job(client: YarnClient, appId: ApplicationId) {
|
|
|
|
|
+
|
|
|
|
|
+ def waitForFinish(timeoutMs: Long): Option[ApplicationStatus] = {
|
|
|
|
|
+ val startTimeMs = System.currentTimeMillis()
|
|
|
|
|
+
|
|
|
|
|
+ while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
|
|
|
|
|
+ val status = getStatus()
|
|
|
|
|
+ status match {
|
|
|
|
|
+ case SuccessfulFinish() | UnsuccessfulFinish() => {
|
|
|
|
|
+ return Some(status)
|
|
|
|
|
+ }
|
|
|
|
|
+ case _ =>
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ Thread.sleep(1000)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ None
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def waitForStatus(status: ApplicationStatus, timeoutMs: Long): Option[ApplicationStatus] = {
|
|
|
|
|
+ val startTimeMs = System.currentTimeMillis()
|
|
|
|
|
+
|
|
|
|
|
+ while (System.currentTimeMillis() - startTimeMs < timeoutMs) {
|
|
|
|
|
+ if (getStatus() == status) {
|
|
|
|
|
+ return Some(status)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ Thread.sleep(1000)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ None
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private def getStatus(): ApplicationStatus = {
|
|
|
|
|
+ val statusResponse = client.getApplicationReport(appId)
|
|
|
|
|
+ convertState(statusResponse.getYarnApplicationState, statusResponse.getFinalApplicationStatus)
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private def convertState(state: YarnApplicationState, status: FinalApplicationStatus): ApplicationStatus = {
|
|
|
|
|
+ (state, status) match {
|
|
|
|
|
+ case (YarnApplicationState.FINISHED, FinalApplicationStatus.SUCCEEDED) => SuccessfulFinish()
|
|
|
|
|
+ case (YarnApplicationState.KILLED, _) | (YarnApplicationState.FAILED, _) => UnsuccessfulFinish()
|
|
|
|
|
+ case (YarnApplicationState.NEW, _) | (YarnApplicationState.SUBMITTED, _) => New()
|
|
|
|
|
+ case _ => Running()
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+trait ApplicationStatus
|
|
|
|
|
+case class New() extends ApplicationStatus
|
|
|
|
|
+case class Running() extends ApplicationStatus
|
|
|
|
|
+case class SuccessfulFinish() extends ApplicationStatus
|
|
|
|
|
+case class UnsuccessfulFinish() extends ApplicationStatus
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
|
|
|
/*
|
|
/*
|
|
|
object Client {
|
|
object Client {
|