【问题标题】:Implementation of Streaming Join in FlinkFlink 中 Streaming Join 的实现
【发布时间】:2021-05-14 21:35:14
【问题描述】:

我正在研究 Flink 中 join 的各种实现。在批处理模式下,我遇到了hybrid-hash join 和sort-merge join。在这两种情况下,在连接之前都会进行阻塞洗牌,因此连接之前的运算符的输出会被具体化到一些非临时存储,如here 所说。

我现在正在查看流连接案例。我已经看到了一个实现,其中为两个输入制作了两个哈希表。每当输入出现时,它都会保存在其哈希表中,并针对其他哈希表进行探测以产生结果。为了限制哈希表的大小,我们在哈希表中放置了一个输入保存的窗口。我的第一个问题是:

Do all stream join cases have this requirement of a window?

具体来说,我想讨论一个大型静态customers 表与Orders 流连接的连接实现。在我看来,物理实现应该是这样的:

customers 表首先进行了哈希分区。然后orders流开始流入。由于执行模式为streaming,所以直接发送orders表加入任务,无需任何物化。

flink 有没有这样的 join 或者我可以在 Flink 中实现这个吗?

【问题讨论】:

    标签: inner-join apache-flink flink-streaming


    【解决方案1】:

    嗯,这正是 BATCH 的实现方式。

    在STREAMING 中,您没有完整的客户表,因为根据定义它是无限的。

    对于BATCH,我只引用这个post from their official blog:

    Flink 具有用于许多操作的流式运行时运算符,但也有用于有界输入的专用运算符 [...] 批处理连接可以将一个输入完全读取到哈希表中,然后使用另一个输入进行探测。流连接需要为双方建表,因为它需要不断地处理两个输入

    此链接还包含有关输入大小的信息:它可以溢出到磁盘。不需要窗口化(但如果您指定它,它肯定会帮助您保持性能/部署大小要求)


    现在,如果您处于 STREAMING 模式并且知道一侧不会改变,您仍然可以将其告知 Flink,以便它围绕它进行优化。 Use JOIN <table> FOR SYSTEM_TIME AS OF <table>.{ proctime | rowtime } for that effect:

    时间连接采用任意表(左侧输入/探测站点)并将每一行与版本化表中对应行的相关版本相关联(右侧输入/构建端)

    但是请注意,如果您使用 JDBC,这些探测端请求将直接通过 Flink 并在数据库中查找(确保您在连接键上有索引)

    【讨论】:

    • 感谢您的回答。时间连接有点类似于我描述的情况。任何想法,如果静态表很大,是否可以使用混合哈希连接在 Flink 中实现这一点?具体来说,我想知道的是,这个连接是否可以分两个阶段实现(或者是否实现) - 首先是阻塞构建分区,然后是流式管道探测。我知道我说得太具体了,但请多多包涵。
    • 好吧,静态表不能大于可以溢出到磁盘的大小。 Hybrid-has-join 虽然是一个非常低级的结构,但 Flink 在其执行计划中选择要做什么。更多详情:flink.apache.org/news/2015/03/13/…。如果您想查看内部结构,请在代码库中查找 HashJoinBuildFirstProperties 和 HINT_LOCAL_STRATEGY_HASH_BUILD_FIRST(您可能会强制优化器使用)
    • 感谢您详细介绍。我确实浏览了您提到的链接,但它只谈论DataSet 而不是DataStream。我想我想问的是 hyrbrid hash join where the build table first comes in and then probe phase starts 在流媒体案例中甚至是 Flink 允许的吗?因为如果构建端也是无限的,那么它就不会完成。还是我可以强制 flink 使用 hash-join 即使是流式传输的情况。
    • 另外,重申一下,我想通过在探测表到来之前对其进行分区来解决大构建表问题。让我们假设它在分区后适合内存。
    • 那么 Flink 的优化器应该可以毫无问题地提出正确的计划。分区/改组将已经由运行时决定。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-01-07
    • 1970-01-01
    • 2020-02-13
    • 2018-06-02
    • 1970-01-01
    相关资源
    最近更新 更多