【问题标题】:stream processing architecture: future events effect past results流处理架构:未来事件影响过去的结果
【发布时间】:2018-01-06 13:41:05
【问题描述】:

我是流处理(kafka 流/flink/storm/spark 等)的新手,并试图找出处理现实世界问题的最佳方法,这里以一个玩具示例为代表。我们在发布订阅/数据摄取方面与 Kafka 相关联,但在流处理器框架/方法方面没有特别的依恋。

理论上,假设我有一个偶尔发出浮点值的源。同样在任何给定点,都有一个乘数 M 应该应用于这个源的值;但是 M 可以改变,而且至关重要的是,我可能会在很久以后才知道这种改变——甚至可能不是“按变更顺序”。

我正在考虑在 Kafka 中将其表示为

"Values": (timestamp, floating point value) - the values from the source, tagged with their emission time.

"Multipliers": (timestamp, floating point multiplier) - indicates M changed to this floating point multiplier at this timestamp.

然后,我很想创建一个输出主题,比如“结果”,使用标准流处理框架连接两个流,并且仅将 Values 中的每个值与乘数确定的当前乘数相乘。

但是,根据我的理解,这是行不通的,因为发布到乘数的新事件可能会对已写入结果流的结果产生任意大的影响。从概念上讲,我希望有一个类似于结果流的东西,它是针对 Values 中的所有值发布到 Multipliers 的最后一个事件的最新内容,但可以在进一步的 Values 或 Multipliers 事件进入时“重新计算”。

使用 kafka 和主要流处理器实现/架构此目标有哪些技术?

例子:

最初,

Values = [(1, 2.4), (2, 3.6), (3, 1.0), (5, 2.2)]
Multipliers = [(1, 1.0)]
Results = [(1, 2.4), (2, 3.6), (3, 1.0), (5, 2.2)]

稍后,

Values = [(1, 2.4), (2, 3.6), (3, 1.0), (5, 2.2)]
Multipliers = [(1, 1.0), (4, 2.0)]
Results = [(1, 2.4), (2, 3.6), (3, 1.0), (5, 4.4)]

最后,在将另一个事件发布到 Multipliers 之后(也发出了一个新值):

Values = [(1, 2.4), (2, 3.6), (3, 1.0), (5, 2.2), (7, 5.0)]
Multipliers = [(1, 1.0), (4, 2.0), (2, 3.0)]
Results = [(1, 2.4), (2, 10.8), (3, 3.0), (5, 4.4), (7, 10.0)]

【问题讨论】:

  • 在这个程序中,Multiplier 通过键将值乘以。所以你的结果会受到影响。
  • 恕我直言,这是相当广泛的,可以给你一个具体的答案。实际的解决方案将取决于要求:“我们需要对数据做什么”。在提供的示例中,我将存储两个流并在读取时执行操作:即。当需要结果时。但这可能还不够,具体取决于实际场景中的应用需求。
  • 好点 maasg。在我们的例子中,有太多的数据流入以支持推迟计算。另外,我们需要做一个查询,比如“给我所有的结果值和它们的时间戳,其中值在 X 和 Y 之间,据你所知,根据当前关于乘数的信息”;无法在未计算结果的情况下为该查询的结果编制索引。

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


【解决方案1】:

我只熟悉 Spark,为了使其按您描述的那样工作,您希望在收到新的乘数时选择性地“更新”以前的结果,同时将最高索引乘数应用于尚未的新值有一个乘数应用于他们。 AFAIK,Spark 本身不会让您使用流式传输来执行此操作(您需要缓存和更新旧结果,并且您还需要知道哪个是用于新值的乘数),但是您可以编写这样的逻辑来编写您将“结果”主题添加到常规数据库表中,并且当您收到新的乘数时,Values 数据框中的所有后续事件都将使用该值,但您将进行一次检查以查找结果表中是否有值现在需要更新以使用新的乘数并简单地更新 DB 表中的这些值。

您的结果使用者必须能够处理插入和更新。您可以将 Spark 与has a connector 的任何数据库一起使用来实现此目的。

或者,您可以使用SnappyData,它将 Apache Spark 变成一个可变的计算 + 数据平台。使用 Snappy,您可以将 Values 和 Multipliers 作为常规流数据帧,并将 Results 作为数据帧设置,作为 SnappyData 中的复制表。当您处理乘数流中的新条目时,您将更新存储在结果表中的所有结果。这可能是完成您正在尝试做的事情的最简单方法

【讨论】:

    猜你喜欢
    • 2019-04-24
    • 2016-08-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多