【问题标题】:Apache Spark streaming - Timeout long-running batchApache Spark 流式传输 - 超时长时间运行的批处理
【发布时间】:2017-12-12 02:33:21
【问题描述】:

我正在设置 Apache Spark 长时间运行的流式传输作业,以使用 InputDStream 执行(非并行化)流式传输。

我想要实现的是,当队列中的批处理花费太长时间(基于用户定义的超时)时,我希望能够跳过该批处理并完全放弃它 - 并继续执行其余部分.

我无法在 spark API 中或在线找到解决此问题的方法——我研究过使用 StreamingContext awaitTerminationOrTimeout,但这会在超时时杀死整个 StreamingContext,而我想要做的只是跳过/杀死当前批次。

我也考虑过使用 mapWithState,但这似乎不适用于这个用例。最后,我正在考虑设置一个 StreamingListener 并在批处理开始时启动一个计时器,然后在达到某个超时阈值时让批处理停止/跳过/杀死,但似乎仍然没有办法杀死批处理。

谢谢!

【问题讨论】:

  • 很好奇为什么 mapWithState 在这里不适用。就像在批处理上创建会话一样?像这样?
  • 好吧,我不使用 Pair DStreams。从理论上讲,如果我是,我也不清楚 API - 如果我确实在某个键上设置了超时,这会做我想要的(跳过批处理中的工作)吗?
  • 这可能很难实现。侦听器将为您提供监视作业运行时间的方法,但我认为取消它会很困难。我查看了(作业调度程序)[github.com/apache/spark/blob/master/streaming/src/main/scala/…,我看不到一个 API 钩子可以在哪里解除批处理的结果。如果你真的需要这个,恐怕你需要修补代码来实现这样的截止日期取消政策。
  • ps:顺便说一句有趣的问题。

标签: apache-spark timeout streaming spark-streaming dstream


【解决方案1】:

我在 yelp 上看过一些文档,但我自己没有做过。

使用UpdateStateByKey(update_func)mapWithState(stateSpec)

  1. 第一次看到事件并初始化状态时附加超时
  2. 如果过期则删除状态

    def update_function(new_events, current_state):
        if current_state is None:
            current_state = init_state()
            attach_expire_datetime(new_events)
            ......
        if is_expired(current_state):
            return None //current_state drops?
        if new_events:
            apply_business_logic(new_events, current_state)
    

看起来结构化流式水印也会在事件超时时丢弃事件,如果这适用于您的作业/阶段超时丢弃。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-12-16
    • 1970-01-01
    • 1970-01-01
    • 2017-05-02
    • 2018-01-12
    • 2018-01-11
    • 2015-04-23
    • 1970-01-01
    相关资源
    最近更新 更多