【发布时间】:2019-06-27 15:01:53
【问题描述】:
我们正在考虑将 flink 用于一个用例,但不确定 flink 是否适合它。这是我的用例。当事件 e1 到达时,我们需要对其进行处理并发出结果。源和接收器与此讨论无关,但您可以将消息队列服务视为源和接收器。一个事件的整个处理独立于其他事件。因此,在处理事件 e1 时,我们不需要 e2 或任何其他事件。作为处理的一部分,我们需要执行步骤 1、步骤 2、步骤 3、步骤 4,如下图所示。注意 step2 和 step3 应该并行进行。
事件的处理延迟对我们来说至关重要。因此,我需要在该元素的处理完成后立即发出结果,而不是等待某个窗口超时。由于我对 Flink 的了解有限,我只能想到下面的方法
DataStream<Map<String, Object>> step1 = env.addSource(...);
DataStream<Map<String, Object>> step2 = step1.map(...);
DataStream<Map<String, Object>> step3 = step1.map(...);
现在,我如何组合 step2 和 step3 的结果并发出结果?在这个简单的示例中,我只有两个要合并的流,但也可以超过 2 个。我可以做一个流的联合。我可以有一个唯一的事件 ID 来对与特定事件相关的中间步骤的输出进行分组。
DataStream<Map<String, Object>> mergedStream = step1.union(step2).keyBy(...);
但是如何发出结果呢?理想情况下,我想说“只要我从 step2 和 step3 获得特定键的输出就发出结果”,而不是“每 30 毫秒发出一次结果”。后者有两个问题:它可能会发出部分结果并且它有延迟。有没有办法指定前者? 我正在探索 Flink,但如果它能解决我的用例,我愿意考虑其他替代方案。
【问题讨论】:
-
所以你想处理 step1,然后是任意数量的并行步骤,然后处理最后一步,其中最后一步需要每个并行步骤的结果,然后才能发出结果?
-
正确。我将在最后一步合并所有并行步骤的结果。
-
谢谢 - 大卫打败了我,下面有一个很好的答案。
标签: apache-flink flink-streaming