【发布时间】: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