爬虫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

智能推荐

python人脸识别用什么库_基于Python的人脸识别库:离线识别率高达99.38%,无敌-程序员宅基地

文章浏览阅读435次。数据测试库Labeled Faces in the Wild:http://vis-www.cs.umass.edu/lfw/模型提供了一个简单的face_recognition命令行工具让用户通过命令就能直接使用图片文件夹进行人脸识别操作。注意:不管你是为了Python就业还是兴趣爱好,记住:项目开发经验永远是核心,如果你没有2020最新python入门到高级实战视频教程,可以去小编的Pyt..._python人脸识别的库

java htmlparser 使用教程_HtmlParser基础教程-程序员宅基地

文章浏览阅读97次。1、相关资料官方文档:http://htmlparser.sourceforge.net/samples.htmlAPI:http://htmlparser.sourceforge.net/javadoc/index.html其它HTML 解释器:jsoup等。由于HtmlParser自2006年以后就再没更新,目前很多人推荐使用jsoup代替它。2、使用HtmlPaser的关键步骤(1)通过Pa..._java使用htmlparser

文本相似度分析(基于jieba和gensim)-程序员宅基地

文章浏览阅读829次。基础概念本文在进行文本相似度分析过程分为以下几个部分进行,文本分词语料库制作算法训练结果预测分析过程主要用两个包来实现jieba,gensimjieba:主要实现分词过程gensim:进行语料库制作和算法训练结巴(jieba)分词在自然语言处理领域中,分词和提取关键词都是对文本处理时通常要进行的步骤。用Python语言对英文文本进行预处理时可选择NLTK库,中文文本预处..._python gensim模块和jieba模块的区别

ScrollPic.js—简单易用的图片左右滚动插件-程序员宅基地

文章浏览阅读5.8k次。ScrollPic.js对于一些新手来说是一个很好理解运用的图片左右滚动插件,兼容性较好,可以放心大胆的使用。_scrollpic.js

Java开发实例大全提高篇——操作PDF篇-程序员宅基地

文章浏览阅读275次。第4篇 操作PDF篇 第13章 操作PDF文档 13.1 文档和文档属性 实例380 创建PDF文档 public static void main(String[] args) { try { Document document = n..._java开发实例大全pdf百度云

java socket缓冲区大小_socket tcp缓冲区大小的默认值、最大值-程序员宅基地

文章浏览阅读1.6k次。Author:阿冬哥Created:2013-4-17Blog:http://blog.csdn.net/c359719435/Copyright 2013阿冬哥http://blog.csdn.net/c359719435/使用以及转载请注明出处1 设置socket tcp缓冲区大小的疑惑疑惑1:通过setsockopt设置SO_SNDBUF、SO_RCVBUF这连个默认缓冲区的值,再用ge..._java api 调用setsockopt(2)系统调用so_rcvbuf选项来控制它的大小

随便推点

JAVA字体颜色背景颜色设置_java setforeground-程序员宅基地

文章浏览阅读3k次,点赞2次,收藏16次。import java.awt.*;   import java.awt.event.*;   public class love extends Frame   {   public static void main(String [] args)   {   lov..._java setforeground

sqlmap之--os-shell_sqlmap --os-shell-程序员宅基地

文章浏览阅读1w次,点赞9次,收藏34次。1、–os-shell原理使用udf提权获取webshell,也是通过into outfile向服务器写入两个文件,一个是可以直接执行系统命令,一个是进行上传文件。–os-shell的执行条件:dbms为mysql,网站必须是root权限攻击者需要知道网站的绝对路径magic_quotes_gpc = off,php主动转移功能关闭2、环境介绍phpstudy+sqlmap3、探测网站根目录python3 sqlmap.py -u "127.0.0.1/sqli-labs-master_sqlmap --os-shell

完美解决隐藏Listview和RecyclerView去掉滚动条和滑动到边界阴影的方案-程序员宅基地

文章浏览阅读3.8w次,点赞19次,收藏32次。转载请标明出处: 本文出自:【Android_Jerry的博客】一、首先是Listview的属性设置设置滑动到顶部和底部的背景或颜色:android:overScrollFooter="@android:color/transparent"android:overScrollHeader="@android:color/transparent"设置滑动到边缘时无效果模式:android:ove_recyclerview去掉滚动

C 中printf无法输出问题_vc++2010printf函数无法输出-程序员宅基地

文章浏览阅读1w次。在c++编程过程中遇到printf()函数无法输出的问题,但是代码没有问题,使用puts()函数可以正常输出。原因为系统缓冲区问题。有三个解决办法:1.添加换行符printf("XXXXXXX \n");2.输出后手动刷新系统缓冲区fflush(stdout);3.预先设定无缓冲区setvbuf(stdout, NULL, _IONBF, 0);..._vc++2010printf函数无法输出

秋招春招总结,经验分享(计算机专业)_计算机秋招难吗-程序员宅基地

文章浏览阅读2.2k次,点赞5次,收藏18次。我很庆幸在七年前选择了计算机专业,虽然选专业完全是听从了命运的安排,直接滑档到了第五个志愿,但是我还是很感谢命运给我这样的安排,遇到了我的本科导师还有几个很好的老师,遇到了几个很好的朋友,回想起来,真好。也正因为本科学校有保研资格,我通过不懈努力来了我的研究生学校,选择了我喜欢的方向,做着学术研究,对我自己的领域说不上如数家珍,也可以算得上有了深入了解。有了一定的研究成果以及研究项目,转眼之间到了毕业的时候,毕业之前经历了漫长的找工作之旅,这趟旅程里充满了焦虑不安以及后悔。找工作的时候很迷茫,不知道选择哪_计算机秋招难吗

【Java架构师面试题】设计模式面试专题(共35题含答案)_设计模式面试题-程序员宅基地

文章浏览阅读3.8k次,点赞4次,收藏26次。设计模式(DesignPattern)是前辈们对代码开发经验的总结,是解决特定问题的一系列套路。它不是语法规定,而是一套用来提高代码可复用性、可维护性、可读性、稳健性以及安全性的解决方案。本篇为设计模式面试专题,总共收录了35道常见面试题及答案解析,希望能帮到你~1、什么是设计模式?就是经过实践验证的用来解决特定环境下特定问题的解决方案2、设计模式用来干什么?寻找合适的对象决定对象的粒度指定对象的接口描述对象的实现运用复用机制重复使用经过实践验证的正确的,用来解决某一类问题的解决方案_设计模式面试题

推荐文章

热门文章

相关标签