【发布时间】: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