【发布时间】:2016-05-03 21:24:45
【问题描述】:
我想在kafka主题接收和处理完成后结束流处理。停止不应像( awaitTerminationOrTimeout )那样特定于时间。有没有办法在主题耗尽后停止 sparkstreamingcontext。有没有办法将 Dstream[T] 与 T 值进行比较来指示控制流?
【问题讨论】:
标签: scala apache-spark spark-streaming
我想在kafka主题接收和处理完成后结束流处理。停止不应像( awaitTerminationOrTimeout )那样特定于时间。有没有办法在主题耗尽后停止 sparkstreamingcontext。有没有办法将 Dstream[T] 与 T 值进行比较来指示控制流?
【问题讨论】:
标签: scala apache-spark spark-streaming
如果流为空,我有 80% 的把握认为 isEmpty 应该返回 true,headOption 应该在 KafkaMessageStream 上为 None。
【讨论】:
Future({() => while(!kafkaMessageStream.isEmpty) { Thread.sleep(100)} sparkTreamingContext.stop(true) }) 这将每 100 毫秒检查一次是否有任何消息,如果没有则停止。
最好的方法是,在开始读取流之前,获取主题中所有分区的最新偏移量,然后检查接收到的偏移量何时到达那里。如果您想了解如何获取主题的偏移量,请参阅我的 previous answer。
流程最终是:
SimpleConsumer
OffsetRequest 返回
最早和最新的偏移量(见上一个答案)OffsetRequest 中收到的最新消息您已完成【讨论】: