【问题标题】:How to process only last, the most relevant, events (and skip the others when latency grew too fast)?如何只处理最后一个最相关的事件(并在延迟增长过快时跳过其他事件)?
【发布时间】:2016-04-10 16:49:22
【问题描述】:

上下文:处理来自 Kafka 的数据并将结果发送回 Kafka。

问题:每个事件可能需要几秒钟的时间来处理(正在进行改进)。在此期间,事件(和 RDD)确实会累积。不需要处理中间事件(按键),只需处理最后的事件。因此,当一个进程完成时,最好让 Spark Streaming 跳过所有不是当前最后一个事件(按键)。

我不确定该解决方案是否可以仅使用 Spark Streaming API 完成。据我了解Spark Streaming,DStream RDD会一一积累处理,以后有没有其他的就不考虑了。

可能的解决方案:

  • 仅使用 Spark Streaming API,但我不确定如何使用。 updateStateByKey 似乎是一个解决方案。但是我不确定当 DStream RDD 累积时它是否会正常工作,并且您只需按键处理 lasts 事件。

  • 有两个 Spark Streaming 管道。一种通过键获取最后更新的事件,将其存储在地图或数据库中。第二个管道仅在事件是另一个管道所指示的最后一个事件时才处理事件。子问题:

    • 两个管道是否可以共享相同的sparkStreamingContext 并以不同的速度(低处理与高处理)处理相同的 DStream?

    • 是否可以在不使用外部数据库的情况下轻松地在管道之间共享值(例如地图)?我认为累加器/广播可以工作,但我不确定在两条管道之间。

【问题讨论】:

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


    【解决方案1】:

    考虑到流媒体是一个连续的过程,很难在这种情况下定义“最后”的含义。但是,假设您要在给定的时间段内处理最后一个事件,例如每 10 秒运行一次处理,并且在这 10 秒帧中只为每个键获取最后一个事件 - 有几种可能的方法。

    窗口方法

    其中一个选项是在DStream 上制作窗口

    val windowStream = dStream.window(Seconds(10), Seconds(10))
    windowStream.forEachRDD { /* process only latest events */ }
    

    在这种情况下,windowStream 将具有 RDD,它在过去 10 秒内组合了所有 RDD 中的键/值,您可以在 forEachRDD 中访问所有它们,就像您最初在单个 RDD 中一样。缺点是它不会提供有关事件如何进入流的事件排序的任何信息,但您可能在值中有事件时间信息或重用来自 Kafka 的偏移量

    updateStateByKey 方法

    基本上按照您的建议 - 它可以让您积累价值。 Databricks 有一个很好的例子来说明如何做到这一点here

    虽然他们在示例中进行了累积,但您可以改为更新键的值

    Kafka 日志压缩

    虽然这并不能取代在 Spark 端处理它的需要,但如果您在 Kafka 中保留事件一段时间,您可能需要考虑使用 Kafka 的 Log Compaction 它不能保证重复项不会从 Kafka 进入 Spark 流,但会通过仅在日志尾部保留最新的键来减少 Kafka 中存储的事件数量。

    【讨论】:

    • 感谢您的回答。我认为我的问题具有误导性。主要问题是当 DStream 中存在延迟和 RDD 队列时,如何更轻松地跳过它们。我需要通过某个键查看最后一个事件,只需对用户 ID 对应于最后一个事件的事件进行长时间处理,然后跳过另一个。
    • Spark 流式传输不会自动为您执行此操作,但 updateStateByKey 最接近您想要的。如果您实现更新功能以仅存储最新事件 - 它会成功。
    猜你喜欢
    • 2020-07-19
    • 1970-01-01
    • 1970-01-01
    • 2018-07-03
    • 1970-01-01
    • 2021-09-25
    • 1970-01-01
    • 1970-01-01
    • 2017-03-27
    相关资源
    最近更新 更多