【问题标题】:Can we use Spark streaming for time based events我们可以将 Spark 流用于基于时间的事件吗
【发布时间】:2019-06-01 07:31:24
【问题描述】:

我有如下要求

  1. 有多个设备根据设备配置生成数据。例如,有两个设备以各自的时间间隔生成数据,假设 d1 每 15 分钟生成一次,d2 每 30 分钟生成一次
  2. 所有这些数据都将发送到 Kafka
  3. 我需要使用数据并根据当前小时生成的值和下一小时生成的第一个值对每个设备执行计算。例如,如果 d1 从上午 12:00 到凌晨 1:00 每 15 分钟生成一次数据,则计算基于该小时生成的值和从上午 1:00 到上午 2:00 生成的第一个值。如果该值不是从凌晨 1:00 到凌晨 2:00 产生的,那么我需要考虑从凌晨 12:00 到凌晨 1:00 的数据并将其保存为数据存储库(时间序列)
  4. 这样会有“n”个设备,每个设备都有自己的配置。在上述场景中,设备 d1 和 d2 每 1 小时生成一次数据。可能还有其他设备每 3 小时、6 小时生成一次数据。

目前这个要求是在 Java 中完成的。由于设备随着计算量的增加而增加,我想知道Spark/Spark Streaming是否可以应用于这种场景?任何关于这类需求的文章都可以分享,这样会有很大的帮助。

【问题讨论】:

    标签: java apache-spark bigdata spark-streaming


    【解决方案1】:

    如果,这是一个很大的如果,计算将是设备方面的,您可以利用主题分区并根据设备数量扩展分区数量。消息按每个分区的顺序传递,这是您需要了解的最强大的概念。

    但是,请注意一些话:

    • 主题数量可能会增加,如果要减少,可能需要清除主题并重新开始。
    • 为了确保设备均匀分布,您可以考虑为每个设备分配一个 guid。
    • 如果计算不涉及某种机器学习库并且可以在纯 java 中完成,那么最好使用普通的旧消费者(或流),而不是通过 Spark-Streaming 抽象它们。级别越低,灵活性越大。

    你可以检查一下。 https://www.confluent.io/blog/how-choose-number-topics-partitions-kafka-cluster

    【讨论】:

    • 是的,但是有 10 万台设备,因此很难维护这么多主题,最终目标是计算过去一小时、过去三小时等设备的值并将它们分桶(范围从 0-100%)
    • 我不建议为每个分区添加一个主题,事实上这将是一个坏主意。分区的数量应该或多或少地用于计算的 CPU 数量。如果您有 4 台服务器,每台服务器有 8 个核心,则您可以使用 32 个分区。这个数字大于 64 应该有一个很好的理由。您可能需要某种缓冲区来存储过去三个小时的数据,或者作为 Redis 的缓存
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-09-09
    • 2013-10-24
    • 1970-01-01
    • 1970-01-01
    • 2011-12-28
    • 2013-02-19
    • 1970-01-01
    相关资源
    最近更新 更多