Scala - Apache YARN 상태 체크

김응진·2022년 12월 14일

Apache Zeppelin에서 Scala로 YARN 클러스터 정보를 API로 받아와서 파싱하는 작업을 했다.

HTTP GET 요청 함수

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
}

YARN API 파싱

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

메모리 사용량 체크 후 알림

profile
Developer

0개의 댓글