【问题标题】:How to implement custom job listener/tracker in Spark?如何在 Spark 中实现自定义作业监听器/跟踪器?
【发布时间】:2017-12-23 03:14:21
【问题描述】:

我有一个像下面这样的类,当我通过命令行运行它时,我想查看进度状态。比如,

10% completed... 
30% completed... 
100% completed...Job done!

我在 yarn 上使用 spark 1.0 并使用 Java API。

public class MyJavaWordCount {
    public static void main(String[] args) throws Exception {
        if (args.length < 2) {
            System.err.println("Usage: MyJavaWordCount <master> <file>");
            System.exit(1);
        }
        System.out.println("args[0]: <master>="+args[0]);
        System.out.println("args[1]: <file>="+args[1]);

        JavaSparkContext ctx = new JavaSparkContext(
                args[0],
                "MyJavaWordCount",
                System.getenv("SPARK_HOME"),
                System.getenv("SPARK_EXAMPLES_JAR"));
        JavaRDD<String> lines = ctx.textFile(args[1], 1);

//      output                                            input   output         
        JavaRDD<String> words = lines.flatMap(new FlatMapFunction<String, String>() {
            //              output       input 
            public Iterable<String> call(String s) {
                return Arrays.asList(s.split(" "));
            }
        });

//          K       V                                                input   K       V 
        JavaPairRDD<String, Integer> ones = words.mapToPair(new PairFunction<String, String, Integer>() {
            //            K       V             input 
            public Tuple2<String, Integer> call(String s) {
                //                K       V 
                return new Tuple2<String, Integer>(s, 1);
            }
        });

        JavaPairRDD<String, Integer> counts = ones.reduceByKey(new Function2<Integer, Integer, Integer>() {
            public Integer call(Integer i1, Integer i2) {
                return i1 + i2;
            }
        });

        List<Tuple2<String, Integer>> output = counts.collect();
        for (Tuple2 tuple : output) {
            System.out.println(tuple._1 + ": " + tuple._2);
        }
        System.exit(0);
    }
}

【问题讨论】:

    标签: java apache-spark


    【解决方案1】:

    如果您使用的是 scala-spark,此代码将帮助您添加 spark 侦听器。

    创建你的 SparkContext

    val sc=new SparkContext(sparkConf) 
    

    现在您可以在 spark 上下文中添加您的 spark 侦听器

    sc.addSparkListener(new SparkListener() {
      override def onApplicationStart(applicationStart: SparkListenerApplicationStart) {
        println("Spark ApplicationStart: " + applicationStart.appName);
      }
    
      override def onApplicationEnd(applicationEnd: SparkListenerApplicationEnd) {
        println("Spark ApplicationEnd: " + applicationEnd.time);
      }
    
    });
    

    Here is 用于监听 Spark 调度中事件的接口列表。

    【讨论】:

      【解决方案2】:

      你应该实现SparkListener。只需覆盖您感兴趣的任何事件(作业/阶段/任务开始/结束事件),然后调用sc.addSparkListener(myListener)

      它不会为您提供直接的基于百分比的进度跟踪器,但至少您可以跟踪正在取得的进度及其粗略率。困难来自于 Spark 阶段的数量是多么不可预测,以及每个阶段的运行时间如何有很大的不同。一个阶段内的进展应该更可预测。

      【讨论】:

        【解决方案3】:

        首先,如果您想跟踪进度,那么您可以考虑spark.ui.showConsoleProgress 请查看@Yijie Shens 的答案(Spark output: log-style vs progress-style)。

        我认为没有必要为这样的事情实现 Spark 监听器。除非你非常具体。


        问题:如何在 Spark 中实现自定义作业监听器/跟踪器?

        You can Use SparkListener and intercept SparkListener events.

        在 Spark 框架中实现的经典示例是 HeartBeatReceiver。

        示例: HeartBeatReceiver.scala

        /**
         * Lives in the driver to receive heartbeats from executors..
         */
        private[spark] class HeartbeatReceiver(sc: SparkContext, clock: Clock)
          extends SparkListener with ThreadSafeRpcEndpoint with Logging {
        
          def this(sc: SparkContext) {
            this(sc, new SystemClock)
          }
        
          sc.addSparkListener(this) ...
        

        以下是可用的侦听器事件列表。哪些应用程序/作业事件应该对您有用

        • SparkListenerApplicationStart

        • SparkListenerJobStart

        • SparkListenerStageSubmitted

        • SparkListenerTaskStart

        • SparkListenerTaskGettingResult

        • SparkListenerTaskEnd

        • SparkListenerStageCompleted

        • SparkListenerJobEnd

        • SparkListenerApplicationEnd

        • SparkListenerEnvironmentUpdate

        • 已添加 SparkListenerBlockManager

        • SparkListenerBlockManagerRemoved

        • SparkListenerBlockUpdated

        • SparkListenerUnpersistRDD

        • SparkListenerExecutor 已添加

        • SparkListenerExecutorRemoved

        【讨论】:

          猜你喜欢
          • 2021-12-22
          • 2016-05-12
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2011-08-21
          • 2011-02-02
          • 2010-10-06
          相关资源
          最近更新 更多