【问题标题】:Flink or Spark for incremental data用于增量数据的 Flink 或 Spark
【发布时间】:2017-08-07 08:38:40
【问题描述】:

我没有使用FlinkSpark 的经验,我想将其中一个用于我的用例。我想介绍我的用例,并希望了解这是否可以用其中任何一个来完成,如果他们都可以做到,那么哪一个效果最好。

我有一堆实体A 存储在数据存储中(准确地说是Mongo,但实际上并不重要)。我有一个 Java 应用程序,它可以加载这些实体并在它们上运行一些逻辑以生成某种数据类型的 Stream E(100% 清楚我没有 Es 在任何数据集,我需要在从数据库加载 As 后用 Java 生成它们)

所以我有这样的东西

A1 -> Stream<E>
A2 -> Stream<E>
...
An -> Stream<E>

数据类型E有点像Excel中的一长行,它有一堆列。我需要收集所有的Es 并像在 Excel 中那样运行某种数据透视聚合。我可以在SparkFlink 中看到如何轻松做到这一点。

现在是我想不通的部分。

假设实体A1 之一被更改(由用户或进程),这意味着A1 的所有Es 都需要更新。当然我可以重新加载我所有的As,重新计算所有Es,然后重新运行整个聚合。我想知道这里是否可以更聪明一点。

是否可以只为A1 重新计算Es 并进行最少的处理。

对于Spark,是否可以保留RDD,并且只在需要时更新其中的一部分(这里是A1Es)?

对于Flink,在流式传输的情况下,是否可以更新已经处理过的数据点?能处理这种情况吗?或者我可以为A1 的旧Es 生成否定 事件(即从结果中删除它们)然后添加新事件?

这是一个常见的用例吗?这甚至是FlinkSpark 的设计目的吗?我会这么想,但我也没有使用过,所以我的理解非常有限。

【问题讨论】:

    标签: java apache-spark batch-processing apache-flink stream-processing


    【解决方案1】:

    我认为您的问题非常广泛,取决于许多条件。在 flink 中,您可以拥有一个 MapState&lt;A, E&gt; 并仅更新更改后的 A's 的值,然后根据您的用例在下游生成更新后的 E's 或生成差异(撤回流)。

    在 Flink 中存在 Dynamics TablesRetraction Streams 的概念,它们可能会激发您的灵感,或者 Table API 可能已经涵盖了您的用例。你可以查看文档here

    【讨论】:

    • 我保持问题的广泛性试图抓住问题的本质。你说答案取决于许多条件。您希望我就问题的哪些方面进行扩展?
    猜你喜欢
    • 2017-01-02
    • 1970-01-01
    • 1970-01-01
    • 2021-12-06
    • 1970-01-01
    • 2019-11-11
    • 2018-09-02
    • 1970-01-01
    • 2023-01-19
    相关资源
    最近更新 更多