【问题标题】:Flink join streams using a disposable one-time join key使用一次性一次性加入密钥的 Flink 加入流
【发布时间】:2018-08-01 04:10:22
【问题描述】:

我有一个关于在 Flink 上加入两个流的问题。我使用了两个不同的数据流,在某些时候我需要 加入他们。每个数据流都被标记了一个唯一的 ID,作为这些流之间的连接点。 没有窗口的概念,所以为了连接这两个数据流,我做了 first.connect(second).keyBy(0,0)。

这似乎有效,因为我得到了正确的结果,但我的担忧是长期的。我没有明确保留任何 执行连接的操作员(coFlatMap)上的状态,但是如果假设一个流(例如第一个)提供唯一的会发生什么 id 和第二个未能提供加入 id(我想对于那些已经加入的操作员丢弃任何类型的内部状态)?内存/状态足迹是不断增长还是存在某种过期机制?

如果是这种情况,我该如何解决这个问题?或者你能建议我另一种方法吗?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    有几种方法可以实现这种连接。

    1. 使用CoProcessFunction。当密钥的第一条记录到达时,您将其存储在 state 中并注册一个计时器,该计时器在 x 分钟/小时/天后触发。当第二条记录到达时,您执行连接并清除状态。如果第二条记录没有到达,则在计时器触发时调用onTimer() 方法。此时,您既可以清除状态并返回(INNER JOIN 语义),也可以转发用null 值填充的第一条记录(OUTER JOIN 语义),清除状态并返回。计时器充当安全网,可以在某个时候移除状态。这取决于您希望等待第二条记录到达多长时间。

    2. Table API 或 SQL 提供了一个时间窗口连接(Table APISQL),其工作方式与我在 1 中描述的类似。不同之处在于窗口连接实现将尝试连接所有记录(即,来自每个输入流的不止一个)在连接间隔期间到达,因此会使状态保持更长时间。但是,一旦时间超过了加入间隔,它就会清除状态。

    3. Flink 1.6.0(将于 2018 年 8 月上旬发布)将包含一个用于 DataStream API 的 interval join,其工作方式类似于 Table API 的窗口连接(逻辑相似,名称不同)。它还将使状态保持比自定义实现更长的时间,该自定义实现基于每个键在每一侧仅出现一次的假设。

    我会选择方法 1。因为它更节省内存并且仍然相当容易实现。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-03-24
      • 2013-10-15
      • 1970-01-01
      • 1970-01-01
      • 2016-11-01
      • 2018-05-20
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多