爬虫Spark UI(Spark streaming监控)_spark爬虫-程序员宅基地

技术标签: spark streaming  爬虫  spark  监控  Spark Streaming  

spark streaming作为实时任务,出了问题并不像离线任务重跑就可以了.对监控要求个方面要求较高.在任务失败 堵塞 卡死等情况下都需要发邮件或者短信报警.
比较普遍的方式是利用spark streaming自带的StreamingListener接口来监控.
如果前者不满足要求,我们也可以简单写个静态爬虫轮询爬取spark ui上的各种指标来diy监控.

方案一 StreamingListener接口

StreamingListener接口只需要新建一个监控类继承StreamingListener,然后重写需要的方法即可.然后记得在主类里加上 ssc.addStreamingListener执行.
以下是个简单的示例,在batch开始时监控schedulingDelay

class StreamingMonitor(ssc:StreamingContext) extends StreamingListener{
  override def onBatchStarted(batchStarted: StreamingListenerBatchStarted): Unit = {
    val Delay_ts = batchStarted.batchInfo.schedulingDelay.get
    if(Delay_ts > DELAY_MAX ){
        sendEmail(...)
    }
  }
}

...
//在main里加
    ssc.addStreamingListener(new StreamingMonitor(ssc))

值得注意的是,StreamingListener接口有多个方法可以重写

//需要监听spark streaming中各个阶段的事件只需实现这个特质中对应的事件函数即可
//本身既有注释说明
trait StreamingListener {

 /** Called when the streaming has been started */
 /** streaming 启动的事件 */
 def onStreamingStarted(streamingStarted: StreamingListenerStreamingStarted) { }

 /** Called when a receiver has been started */
 /** 接收启动事件 */
 def onReceiverStarted(receiverStarted: StreamingListenerReceiverStarted) { }

 /** Called when a receiver has reported an error */
 def onReceiverError(receiverError: StreamingListenerReceiverError) { }

 /** Called when a receiver has been stopped */
 def onReceiverStopped(receiverStopped: StreamingListenerReceiverStopped) { }

 /** Called when a batch of jobs has been submitted for processing. */
 /** 每个批次提交的事件 */
 def onBatchSubmitted(batchSubmitted: StreamingListenerBatchSubmitted) { }

 /** Called when processing of a batch of jobs has started.  */
 /** 每个批次启动的事件 */
 def onBatchStarted(batchStarted: StreamingListenerBatchStarted) { }

 /** Called when processing of a batch of jobs has completed. */
 /** 每个批次完成的事件  */
 def onBatchCompleted(batchCompleted: StreamingListenerBatchCompleted) { }

 /** Called when processing of a job of a batch has started. */
 def onOutputOperationStarted(
     outputOperationStarted: StreamingListenerOutputOperationStarted) { }

 /** Called when processing of a job of a batch has completed. */
 def onOutputOperationCompleted(
     outputOperationCompleted: StreamingListenerOutputOperationCompleted) { }
}

独立监控程序监控spark ui

有些指标可能我们用接口实现不了. 使用scala实现一个简单的静态爬虫来监控spark ui,相当于一个脚本程序替我们不停查看spark ui

//需要jsoup包来读接口并解析html
        <dependency>
            <groupId>org.jsoup</groupId>
            <artifactId>jsoup</artifactId>
        </dependency>

以下是个简单的示例来监控job页面的job执行时间是否过长.
注意有些公司 spark ui可能需要cookie信息.


object StreamingUIMonitorJob {

    val app_id = args(0)
    var dc = getJobDoc(app_id)
    if(dc == null){
      println("get Job document failed")    
      ...
    }

    //如果没有active job则等待1min,等待10min报警
    var active_job_table = dc.getElementById("activeJob-table")
    var alarm_num = 0
    while (active_job_table == null){
      if(alarm_num > 10){
        ...
      }
      Thread.sleep(60000)
      dc = getJobDoc(app_id)
      active_job_table = dc.getElementById("activeJob-table")
      alarm_num = alarm_num + 1
    }


    var durs = Array[String]()
    var batchs = Array[String]()

    //只有一个tbody tbody里可能有多个tr,一个job 一个tr
    val active_jobs = active_job_table.getElementsByTag("tbody")(0).getElementsByTag("tr")
    for(active_job <- active_jobs){
      val job_infos = active_job.getElementsByTag("td")
      println(job_infos(2).text()) //batch 时间
      batchs :+= job_infos(2).text()
      println(job_infos(3).text()) //dur 时间
      durs :+= job_infos(3).text()
    }


    if(durs.length > 0){
      try{
        var delay:Double = 0.0
        for(dur <- durs){
          val ls = dur.split(" ")
          val num = ls(0).toDouble
          if(ls(1) == "min" && num > delay) delay = num
        }
        if(delay > MAX_DELAY_TIME) {
          println(s"active job delay ${delay} min !")
          //报警
          ...
        }
      }catch {
        case e:Exception=>
          println(s"解析durs出错 Exception:${e}")
      }
    }
}

  def getJobDoc(app_id :String): Document ={
    var con = Jsoup.connect(SPARK_UI_URL + app_id )
    var cookie = ""
    // get cookie
    try {
      cookie = con.execute().cookies().toString
      cookie = cookie.substring(cookie.indexOf("{")+1,cookie.lastIndexOf("}"))
    }catch {
      case e:Exception=>
        println("get cookie exception!")
    }
    // get document
    var doc :Document = null
    try {
      con = Jsoup.connect(SPARK_UI_URL + app_id + "/jobs/?proxyapproved=true").header("Cookie",cookie)
      doc = con.get()
    }catch {
      case e:Exception=>
        println("get Job document exception!")
    }
    doc
  }

这只是一种场景需要,实际上spark ui上的所有指标都可以通过jsoup解析html监控的.

版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/qq_21277411/article/details/102955107

智能推荐

while循环&CPU占用率高问题深入分析与解决方案_main函数使用while(1)循环cpu占用99-程序员宅基地

文章浏览阅读3.8k次,点赞9次,收藏28次。直接上一个工作中碰到的问题,另外一个系统开启多线程调用我这边的接口,然后我这边会开启多线程批量查询第三方接口并且返回给调用方。使用的是两三年前别人遗留下来的方法,放到线上后发现确实是可以正常取到结果,但是一旦调用,CPU占用就直接100%(部署环境是win server服务器)。因此查看了下相关的老代码并使用JProfiler查看发现是在某个while循环的时候有问题。具体项目代码就不贴了,类似于下面这段代码。​​​​​​while(flag) {//your code;}这里的flag._main函数使用while(1)循环cpu占用99

【无标题】jetbrains idea shift f6不生效_idea shift +f6快捷键不生效-程序员宅基地

文章浏览阅读347次。idea shift f6 快捷键无效_idea shift +f6快捷键不生效

node.js学习笔记之Node中的核心模块_node模块中有很多核心模块,以下不属于核心模块,使用时需下载的是-程序员宅基地

文章浏览阅读135次。Ecmacript 中没有DOM 和 BOM核心模块Node为JavaScript提供了很多服务器级别,这些API绝大多数都被包装到了一个具名和核心模块中了,例如文件操作的 fs 核心模块 ,http服务构建的http 模块 path 路径操作模块 os 操作系统信息模块// 用来获取机器信息的var os = require('os')// 用来操作路径的var path = require('path')// 获取当前机器的 CPU 信息console.log(os.cpus._node模块中有很多核心模块,以下不属于核心模块,使用时需下载的是

数学建模【SPSS 下载-安装、方差分析与回归分析的SPSS实现(软件概述、方差分析、回归分析)】_化工数学模型数据回归软件-程序员宅基地

文章浏览阅读10w+次,点赞435次,收藏3.4k次。SPSS 22 下载安装过程7.6 方差分析与回归分析的SPSS实现7.6.1 SPSS软件概述1 SPSS版本与安装2 SPSS界面3 SPSS特点4 SPSS数据7.6.2 SPSS与方差分析1 单因素方差分析2 双因素方差分析7.6.3 SPSS与回归分析SPSS回归分析过程牙膏价格问题的回归分析_化工数学模型数据回归软件

利用hutool实现邮件发送功能_hutool发送邮件-程序员宅基地

文章浏览阅读7.5k次。如何利用hutool工具包实现邮件发送功能呢?1、首先引入hutool依赖<dependency> <groupId>cn.hutool</groupId> <artifactId>hutool-all</artifactId> <version>5.7.19</version></dependency>2、编写邮件发送工具类package com.pc.c..._hutool发送邮件

docker安装elasticsearch,elasticsearch-head,kibana,ik分词器_docker安装kibana连接elasticsearch并且elasticsearch有密码-程序员宅基地

文章浏览阅读867次,点赞2次,收藏2次。docker安装elasticsearch,elasticsearch-head,kibana,ik分词器安装方式基本有两种,一种是pull的方式,一种是Dockerfile的方式,由于pull的方式pull下来后还需配置许多东西且不便于复用,个人比较喜欢使用Dockerfile的方式所有docker支持的镜像基本都在https://hub.docker.com/docker的官网上能找到合..._docker安装kibana连接elasticsearch并且elasticsearch有密码

随便推点

Python 攻克移动开发失败!_beeware-程序员宅基地

文章浏览阅读1.3w次,点赞57次,收藏92次。整理 | 郑丽媛出品 | CSDN(ID:CSDNnews)近年来,随着机器学习的兴起,有一门编程语言逐渐变得火热——Python。得益于其针对机器学习提供了大量开源框架和第三方模块,内置..._beeware

Swift4.0_Timer 的基本使用_swift timer 暂停-程序员宅基地

文章浏览阅读7.9k次。//// ViewController.swift// Day_10_Timer//// Created by dongqiangfei on 2018/10/15.// Copyright 2018年 飞飞. All rights reserved.//import UIKitclass ViewController: UIViewController { ..._swift timer 暂停

元素三大等待-程序员宅基地

文章浏览阅读986次,点赞2次,收藏2次。1.硬性等待让当前线程暂停执行,应用场景:代码执行速度太快了,但是UI元素没有立马加载出来,造成两者不同步,这时候就可以让代码等待一下,再去执行找元素的动作线程休眠,强制等待 Thread.sleep(long mills)package com.example.demo;import org.junit.jupiter.api.Test;import org.openqa.selenium.By;import org.openqa.selenium.firefox.Firefox.._元素三大等待

Java软件工程师职位分析_java岗位分析-程序员宅基地

文章浏览阅读3k次,点赞4次,收藏14次。Java软件工程师职位分析_java岗位分析

Java:Unreachable code的解决方法_java unreachable code-程序员宅基地

文章浏览阅读2k次。Java:Unreachable code的解决方法_java unreachable code

标签data-*自定义属性值和根据data属性值查找对应标签_如何根据data-*属性获取对应的标签对象-程序员宅基地

文章浏览阅读1w次。1、html中设置标签data-*的值 标题 11111 222222、点击获取当前标签的data-url的值$('dd').on('click', function() { var urlVal = $(this).data('ur_如何根据data-*属性获取对应的标签对象

推荐文章

热门文章

相关标签