【发布时间】:2018-08-01 04:10:22
【问题描述】:
我有一个关于在 Flink 上加入两个流的问题。我使用了两个不同的数据流,在某些时候我需要 加入他们。每个数据流都被标记了一个唯一的 ID,作为这些流之间的连接点。 没有窗口的概念,所以为了连接这两个数据流,我做了 first.connect(second).keyBy(0,0)。
这似乎有效,因为我得到了正确的结果,但我的担忧是长期的。我没有明确保留任何 执行连接的操作员(coFlatMap)上的状态,但是如果假设一个流(例如第一个)提供唯一的会发生什么 id 和第二个未能提供加入 id(我想对于那些已经加入的操作员丢弃任何类型的内部状态)?内存/状态足迹是不断增长还是存在某种过期机制?
如果是这种情况,我该如何解决这个问题?或者你能建议我另一种方法吗?
【问题讨论】:
标签: apache-flink flink-streaming