更新时间:2026-06-29 GMT+08:00
分享

FlinkSQL支持JTL(Join-To-Live)

Flink双流Join场景中,需要将数据存储在状态后端,目前普遍使用RocksDB作为状态后端,并依赖TTL(Time-To-Live,数据存活时间)机制来管理状态数据的生命周期。在TTL过大或无法确定TTL或数据流量增加的场景下,大流量会导致状态数据增加,增加存储压力,从而导致作业稳定性下降,或TTL过期可能出现数据关联不准确的问题

对于数据关联次数确定的业务,可以使用JTL(Join-To-Live)特性,JTL指根据关联次数确定数据是否过期,当数据达到设定的关联次数阈值后即从状态中清除,从而减少状态后端压力其可以通过如下两种方式配置:

  • 仅支持Flink Regular Joins中的inner join语句。
  • 不支持与作业级TTL、表级TTL、小表广播特性同时使用。
  • 使用JTL特性的表需要指定主键,否则可能导致结果不准确。
  • 方式一:可通过SQL Hint使用

    eliminate-state.left.threshold:表示左关联次数的阈值,单位为次当左表中某条数据的关联次数超过阈值后,该条数据就会过期。适用于左表数据关联次数可确定的场景。

    eliminate-state.right.threshold:表示右关联次数的阈值,单位为次当右表中某条数据的关联次数超过阈值后,该条数据就会过期。适用于表数据关联次数可确定的场景。

    示例一:

    SELECT * FROM t1
       JOIN /*+ OPTIONS('eliminate-state.right.threshold'='1', 'eliminate-state.left.threshold'='2') */
       t2 ON a1 = a2

    示例二:

    SELECT a1, a2, a3 from
      t1
      join /*+ OPTIONS('eliminate-state.left.threshold'='1', 'eliminate-state.right.threshold'='2') */
      t2
      on a1 = a2
      join /*+ OPTIONS('eliminate-state.left.threshold'='3', 'eliminate-state.right.threshold'='4') */
      t3
      on a2 = a3
  • 方式二:通过在“客户端安装路径/Flink/flink/conf/flink-conf.yaml”全局配置使用
    table.exec.join.eliminate-state.left.threshold
    table.exec.join.eliminate-state.right.threshold

相关文档