Apache Zeppelin에서 Scala로 YARN 클러스터 정보를 API로 받아와서 파싱하는 작업을 했다.
import scala.util.parsing.json._
import java.text.SimpleDateFormat
import org.joda.time.format.DateTimeFormat
import org.joda.time.DateTime
import okhttp3._
val fmt_yyyymmdd = DateTimeFormat.forPattern("yyyyMMdd")
var p_date = fmt_yyyymmdd.print(new DateTime())
val fmt_hhmmss = DateTimeFormat.forPattern("HH:mm:ss")
var now_time = fmt_hhmmss.print(new DateTime())
//HTTP GET 요청
def reqHttpGet(url: String): String = {
import java.net.{URL, HttpURLConnection}
import java.io.{DataOutputStream, OutputStreamWriter}
val connection = (new URL(url)).openConnection.asInstanceOf[HttpURLConnection]
connection.setConnectTimeout(30000)
connection.setReadTimeout(30000)
connection.setRequestMethod("GET")
val inputStream = connection.getInputStream
val content = scala.io.Source.fromInputStream(inputStream).mkString
if (inputStream != null) inputStream.close()
content
}
val result = reqHttpGet("http://YARN-RM-IP:8088/ws/v1/cluster/metrics")
val json_result = JSON.parseFull(result)
val get_result = json_result.getOrElse(0).asInstanceOf[Map[String,Any]]
val clusterMetrics = get_result("clusterMetrics").asInstanceOf[Map[Any,Any]]
var appsSubmitted = clusterMetrics("appsSubmitted").asInstanceOf[Double].toInt
var appsPending = clusterMetrics("appsPending").asInstanceOf[Double].toInt
var appsRunning = clusterMetrics("appsRunning").asInstanceOf[Double].toInt
var appsCompleted = clusterMetrics("appsCompleted").asInstanceOf[Double].toInt
var containersAllocated = clusterMetrics("containersAllocated").asInstanceOf[Double].toInt
var availableMB = clusterMetrics("availableMB").asInstanceOf[Double]
var allocatedMB = clusterMetrics("allocatedMB").asInstanceOf[Double]
var totalMB = clusterMetrics("totalMB").asInstanceOf[Double]
var reservedMB = clusterMetrics("reservedMB").asInstanceOf[Double]
var allocatedVirtualCores = clusterMetrics("allocatedVirtualCores").asInstanceOf[Double].toInt
var totalVirtualCores = clusterMetrics("totalVirtualCores").asInstanceOf[Double].toInt
var reservedVirtualCores = clusterMetrics("reservedVirtualCores").asInstanceOf[Double].toInt
var activeNodes = clusterMetrics("activeNodes").asInstanceOf[Double].toInt
var decommissionedNodes = clusterMetrics("decommissionedNodes").asInstanceOf[Double].toInt
var lostNodes = clusterMetrics("lostNodes").asInstanceOf[Double].toInt
var unhealthyNodes = clusterMetrics("unhealthyNodes").asInstanceOf[Double].toInt
var rebootedNodes = clusterMetrics("rebootedNodes").asInstanceOf[Double].toInt
if(availableMB > 1024) availableMB = availableMB / 1024
if(allocatedMB > 1024) allocatedMB = allocatedMB / 1024
if(totalMB > 1024) totalMB = totalMB / 1024
var memory_usage_rate = allocatedMB/totalMB*100