【问题标题】:Create a timer within a spark context在火花上下文中创建计时器
【发布时间】:2016-09-07 00:41:01
【问题描述】:

我们有一个 Spark 应用程序,它从 Kafka 流式传输,并消耗客户活动。我正在尝试确定客户是否在我们的系统上停止活动 3 分钟(即 3 分钟内没有收到该客户的另一笔交易)。

我不确定我是否正在尝试以正确的方式实现这一点,或者在 Spark 中使用此逻辑是否没有意义,但我正在尝试使用 RecurringTimer 类来执行此操作。有没有人实现过类似的东西,如果有的话,spark 库中使用了什么实用程序函数?

任何例子,指针等也将不胜感激

【问题讨论】:

    标签: java apache-spark apache-kafka bigdata


    【解决方案1】:

    看看mapWithState,基本上你会聚合成一个键/值对,由客户的一些标识符和最后收到的交易的时间戳组成。

    执行此聚合后的每个微批处理,您可以检查并查看其中是否有任何用户拥有 timestamp < now() - 3min 并执行某些操作(即,将消息推送到另一个 kafka 队列等)

    mapWithState 上的示例可用here

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-30
    • 1970-01-01
    • 2020-06-09
    • 1970-01-01
    • 2017-07-05
    • 1970-01-01
    相关资源
    最近更新 更多