【问题标题】:stop spark Streaming context kafkaDirectStream停止火花流上下文 kafkaDirectStream
【发布时间】:2016-05-03 21:24:45
【问题描述】:

我想在kafka主题接收和处理完成后结束流处理。停止不应像( awaitTerminationOrTimeout )那样特定于时间。有没有办法在主题耗尽后停止 sparkstreamingcontext。有没有办法将 Dstream[T] 与 T 值进行比较来指示控制流?

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    如果流为空,我有 80% 的把握认为 isEmpty 应该返回 true,headOption 应该在 KafkaMessageStream 上为 None。

    【讨论】:

    • rdd.isEmpty of kafkaStream.foreachRDD 确实返回 true。但它不会使用 sparkStreamingContext.stop(true) 停止应用程序
    • 您可以执行以下操作:Future({() => while(!kafkaMessageStream.isEmpty) { Thread.sleep(100)} sparkTreamingContext.stop(true) }) 这将每 100 毫秒检查一次是否有任何消息,如果没有则停止。
    • 您可能会得到误报,例如,如果代理出现问题的时间长于您的批处理间隔。
    • 我开始使用 sys.exit() 而不是 sparkStreamingContext.stop() 因为后者不会停止驱动程序中的应用程序。
    【解决方案2】:

    最好的方法是,在开始读取流之前,获取主题中所有分区的最新偏移量,然后检查接收到的偏移量何时到达那里。如果您想了解如何获取主题的偏移量,请参阅我的 previous answer。

    流程最终是:

    1. 获取主题的分区和代理
    2. 为每个代理创建一个SimpleConsumer
    3. 对于每个分区,执行 OffsetRequest 返回 最早和最新的偏移量(见上一个答案)
    4. 然后在阅读消息时,检查收到消息的偏移量 相对于分区的已知最后偏移量
    5. 一旦每个分区收到的所有偏移量都与 您的OffsetRequest 中收到的最新消息您已完成

    【讨论】:

      猜你喜欢
      • 2020-05-29
      • 1970-01-01
      • 2016-02-18
      • 2015-12-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-08-04
      • 1970-01-01
      相关资源
      最近更新 更多