【问题标题】:Parallel design of program working with Flink and scala使用 Flink 和 scala 并行设计程序
【发布时间】:2018-01-24 14:52:33
【问题描述】:

这是上下文:

  1. 有一个输入事件流,
  2. 有一些方法可以应用于 流,它应用不同的逻辑来评估每个事件, 说这是“好”或“坏”事件。
  3. 一个事件可以是一个真正的“好”事件只有当它通过所有方法,否则它是一个“坏”事件。
  4. 有一个输出事件流具有事件的结果及其事件ID。

为了解决这个问题,我有两个想法:

  1. 我们可以将每种方法依次应用于每个事件。 但是这是一种批处理,没有应用流处理的优点,同时需要Time(M(ethod)1) + Time(M2) + Time(M3) + .....,可能不适合实时处理。
  2. 我们可以将输入流传递给每个方法,然后我们可以并行运行每个方法,每个方法保存坏事件到永久存储中,然后Main 方法可以查询永久存储以获取每个事件的结果。 但是这有一些问题需要解决:

    • 如何在编程语言(例如 Scala)中并行执行方法,性能如何(网络、CPU、内存)

    • 如何解决同步问题?可以肯定的是,这些方法需要一些时间来计算flag并将其保存到永久存储中,但是Main需要更少的时间来查询flag,这会出现延迟问题。

这不是技术和设计的问题,我想问问你们的想法,如果你有一些新的想法或想法来解决这个问题?期待您的意见。

【问题讨论】:

  • 按顺序执行 (#1)。如果时间成为问题,您始终可以对流进行分区,并并行处理事件。
  • @Dima 有链接让我看懂吗?
  • 不知道你想了解什么

标签: scala parallel-processing apache-flink flink-streaming


【解决方案1】:

并行流,每个流按顺序执行全套评估,是更直接的解决方案。但是,如果这会导致过多的延迟,那么您可以将要并行执行的评估展开,然后将结果重新组合在一起以做出决定。

要进行扇出,请查看 DataStream 上的拆分操作,或使用侧输出。但在执行此 n 路扇出之前,请确保每个事件都有唯一的 ID。如有必要,为每个事件添加一个包含随机数的字段以用作唯一 ID。稍后我们将使用这个唯一 ID 作为键来收集每个事件的所有部分结果。

一旦事件流被拆分,流的每个副本都可以使用 MapFunction 来计算其中一种评估方法。

将给定事件的所有这些单独评估重新收集在一起有点复杂。这里一种合理的方法是将所有结果流联合在一起,然后通过上述唯一 ID 对联合流进行键控。这将汇集每个事件的所有单独结果。然后,您可以使用 RichFlatMapFunction(使用 Flink 的键控、托管状态)在一个地方收集单独评估的结果。一旦给定事件的完整评估集到达这个有状态的平面地图运算符,它就可以计算并发出最终结果。

【讨论】:

    猜你喜欢
    • 2012-02-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-07-23
    • 1970-01-01
    • 2016-05-24
    • 2020-01-17
    • 2013-05-13
    相关资源
    最近更新 更多