【问题标题】:Stream Joins for Large Time Windows with Flink使用 Flink 进行大型时间窗口的流连接
【发布时间】:2019-10-01 16:44:15
【问题描述】:

我需要基于一个键加入两个事件源。事件之间的间隔最长可达 1 年(即 id1 的 event1 可能今天到达,而来自第二个事件源的 id1 对应的 event2 可能会在一年后到达)。假设我只想流式输出连接的事件输出。

我正在探索将 Flink 与 RocksDB 后端一起使用的选项(我遇到了似乎适合我的用例的 Table API)。我无法找到执行这种长窗口连接的参考架构。我预计系统每天可以处理大约 2 亿个事件。

问题:

  1. 使用 Flink 进行这种长窗口连接有什么明显的限制/陷阱吗?

  2. 关于处理这种长窗口连接的任何建议

相关:我也在探索使用 Lambda 和 DynamoDB 作为状态来进行流连接 (Related Question)。如果此信息相关,我将使用托管 AWS 服务。

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    这个用例的明显挑战是一年的大连接窗口大小和可能导致巨大状态大小的高摄取率。

    这里的主要问题是这是否是 1:1 连接,即流 A 中的记录是否与流 B 中的记录完全(或最多)连接一次。这很重要,因为如果您有 1 :1 加入,您可以在将一条记录与另一条记录合并后立即从该州中删除该记录,您无需将其保留一整年。因此,您所在的州只存储尚未加入的记录。假设大多数记录都很快加入,您的状态可能会保持合理的小。

    如果你有一个 1:1 的连接,那么 Flink 的 Table API(和 SQL)的时间窗口连接和 DataStream API 的间隔连接不是你想要的。它们被实现为 m:n 连接,因为每条记录都可能与另一个输入的多个记录连接。因此,它们会在整个窗口间隔内保存 all 记录,即在您的用例中保存一年。如果您有 1:1 连接,您应该自己将连接实现为 KeyedCoProcessFunction

    如果每条记录可以在一年内多次加入,则无法缓冲这些记录。在这种情况下,您可以使用 Flink 的 Table API(和 SQL)的 time-window joins 和 DataStream API 的 Interval join。

    【讨论】:

    • @Fabian Hueske 感谢您的快速回复。我很想知道庞大的国家规模的含义? 1. 是否影响性能(加入性能/吞吐量) 2. 是否影响运营开销(集群维护)(假设 AWS 是我唯一的选择)
    • 流式应用程序需要确保其状态在发生故障时不会丢失。 Flink 执行定期检查点,这意味着它将应用程序状态复制到远程持久存储系统,这对于更大的状态变得更加昂贵。为了恢复,状态被加载回来。同样,如果状态非常大,这会更昂贵。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-18
    • 1970-01-01
    • 1970-01-01
    • 2018-06-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多