【问题标题】:How to continue batch processing after invoking stop() on Spark Streaming Context在 Spark Streaming Context 上调用 stop() 后如何继续批处理
【发布时间】:2018-06-27 13:48:42
【问题描述】:

我们正在撰写一篇关于 Spark 批处理/流处理效率的论文。我们试图检测大量数据中的异常,我们需要的是哪个日志行经历了哪个过程。

因此,我们创建了一个事件模拟,在每个流程之前和之后,我们记录该线路到达/离开该阶段的时间。

但我们面临的一个问题是,我们不希望将分析流处理的时间包含在这些计算中。所以我们基本上需要的是

使用流式进行一些计算, 调用ssc.stop(false,true)(通过HTTP或检测文件结尾), 继续处理有关性能的分析

但是 Spark 的问题是它不允许我们在调用 stop 之后处理 DStream。有什么办法可以复制我们最后一个 DStream 以便我们在调用 stop() 后访问它的对象?

我们在尝试这样做时遇到的错误是:

    Exception in thread "main" java.lang.IllegalStateException: Adding new inputs, transformations, and output operations after stopping a context is not supported

代码架构基本上是这样的:

    val sparkConf = new SparkConf().setAppName("DirectKafkaWordCount").setMaster("local[" + CPUNumber + "]")
    val ssc = new StreamingContext(sparkConf, Seconds(1))

    //Some ml algorithms
    val x = b.map(something)
    ssc.start()

    ssc.awaitTermination()
    ssc.stop(false)

    //Some analytical tracking map reduce jobs
    val y = x.map(getanalytics)

提前感谢,任何想法都非常感谢

【问题讨论】:

    标签: apache-spark apache-kafka streaming spark-streaming


    【解决方案1】:

    你可以这样做:

    b.foreachRDD(r -> stopwatch = Stopwatch.createStarted())
    b.map(something)
    b.foreachRDD(r -> System.out.println(stopwatch)
    ssc.start()
    

    那么你可以在每一轮中消耗时间。

    【讨论】:

    • 是的,但就像我上面解释的那样,我们有大量数据,我们不想根据测量流部分的效率来进行计算。你的建议就是这样做的。
    • 相信我,它在运行时方面发生了很大变化,分析跟踪减少作业甚至比 ML 算法本身花费更多时间。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-06-24
    • 1970-01-01
    • 2020-01-05
    • 1970-01-01
    • 1970-01-01
    • 2019-06-25
    • 2017-12-18
    相关资源
    最近更新 更多