【问题标题】:What are the options to process timeseries data from a Kinesis stream处理 Kinesis 流中的时间序列数据的选项有哪些
【发布时间】:2016-05-10 09:03:44
【问题描述】:

我需要处理来自 AWS Kinesis 流的数据,该流从设备收集事件。在过去 10 秒内收到的所有事件必须每秒调用处理函数。


假设,我有两个设备 A 和 B 将事件写入流。 我的过程名称为 MyFunction 并采用以下参数:

  • 设备ID
  • 一段时间的数据数组

如果我在 10:00:00 开始处理(并且在过去 10 秒内已经积累了设备 A 和 B 的事件) 那么我需要打两个电话:

  • MyFunction(А, {Events for device A from 09:59:50 to 10:00:00})
  • MyFunction(B, {Events for device B from 09:59:50 to 10:00:00})

下一秒,10:00:01

  • MyFunction(А, {Events for device A from 09:59:51 to 10:00:01})
  • MyFunction(B, {Events for device B from 09:59:51 to 10:00:01})

等等。


看起来从设备接收到的所有数据最简单的方法就是将其存储在临时缓冲区中(当然只有最后 10 秒),所以我想先尝试一下。

我发现保存这种基于内存的缓冲区的最方便的方法是创建一个基于 Java Kinesis Client Library (KCL) 的应用程序。

我也考虑过基于 AWS Lambda 的解决方案,但看起来不可能将数据保存在内存中以供 lambda 使用。 Lambda 的另一种选择是拥有 2 个函数,第一个函数必须将所有数据写入 DynamoDB,第二个函数每秒调用一次以处理从 db 而不是从内存中获取的数据。 (所以这个选项要复杂得多)

所以我的问题是:还有哪些其他选项可以实现此类处理?

【问题讨论】:

    标签: amazon-web-services aws-lambda amazon-kinesis amazon-kcl


    【解决方案1】:

    因此,您所做的称为“窗口操作”(或“窗口计算”)。有多种方法可以实现这一点,就像你说的缓冲是最好的选择。

    • 在内存缓存系统中:Ehcache、Hazelcast

    在缓存系统中累积数据并选择适当的驱逐策略(在您的情况下为 10 分钟)。然后进行分组求和运算,计算输出。

    • 内存数据库:Redis、VoltDB

    就像缓存系统一样,您可以使用数据库架构。 Redis 可能是有帮助的和有状态的。如果你使用 VoltDB 或类似的 SQL 系统,调用“sum()”或“avg()”操作会更容易。

    可以使用 Spark 进行计数。您可以尝试 Elastic MapReduce (EMR),这样您就可以留在 AWS 生态系统中并且更容易集成。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-09-23
      • 2011-12-14
      • 1970-01-01
      • 2016-07-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多