多流Join场景支持配置表级别的TTL时间
操作场景
在Flink双流Join场景下,如果Join的左表和右表其中一个表数据变化快,需要较短时间的过期时间,而另一个表数据变化较慢,需要较长时间的过期时间。目前Flink只有表级别的TTL(Time To Live:生存时间),为了保证Join的准确性,需要将表级别的TTL设置为较长时间的过期时间,此时状态后端中保存了大量的已经过期的数据,给状态后端造成了较大的压力。为了减少状态后端的压力,可以单独为左表和右表设置不同的过期时间。其中,左表指JOIN语句中JOIN关键字左侧的表,右表指JOIN语句中JOIN关键字右侧的表。不支持where子句。
/*+ OPTIONS('state.ttl.left'='60S', 'state.ttl.right'='120S') */ 约束与限制
本章节适用于MRS 3.3.0及以后版本。
相关概念
- TTL(Time To Live,生存时间):Flink状态数据的存活时长,超过该时长的状态数据将被自动清理,用于控制状态后端的数据保留周期,防止状态无限增长。
- Hint(SQL提示):一种通过SQL注释语法(/*+ OPTIONS(...) */)向查询引擎传递配置参数的机制,在不修改SQL逻辑的前提下影响查询的执行行为。
- 状态后端(State Backend):Flink用于存储算子状态(如Join缓存数据)的存储引擎。在双流Join场景中,左右表的数据均缓存在状态后端中等待匹配。
- 双流Join:指两个无界数据流(如两个Kafka数据源)之间的Join操作。由于数据流是无界的,Flink需要将已到达的数据缓存在状态后端中,等待对端数据流中匹配的数据到达后执行Join。
在SQL语句中配置示例
- 示例1:为左表和右表分别设置过期时间。
CREATE TABLE user_info (`user_id` VARCHAR, `user_name` VARCHAR) WITH ( 'connector' = 'kafka', 'topic' = 'user_info_001', 'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'latest-offset', 'value.format' = 'csv' ); CREATE table print( `user_id` VARCHAR, `user_name` VARCHAR, `score` INT ) WITH ('connector' = 'print'); CREATE TABLE user_score (user_id VARCHAR, score INT) WITH ( 'connector' = 'kafka', 'topic' = 'user_score_001', 'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'latest-offset', 'value.format' = 'csv' ); INSERT INTO print SELECT t.user_id, t.user_name, d.score FROM user_info as t JOIN -- 为左表和右表设置不同的TTL时间 /*+ OPTIONS('state.ttl.left'='60S', 'state.ttl.right'='120S') */ user_score as d ON t.user_id = d.user_id;
Kafka Broker实例IP地址及端口号说明:- 服务的实例IP地址可通过登录FusionInsight Manager后,单击“集群 > 服务 > Kafka > 实例”,在实例列表页面中查询。
- 集群已启用Kerberos认证(安全模式)时Broker端口为“sasl.port”参数的值,默认为“21007”。
- 集群未启用Kerberos认证(普通模式)时Broker端口为“port”的值,默认为“9092”。如果配置端口号为9092,则需要配置“allow.everyone.if.no.acl.found”参数为true,具体操作如下:
登录FusionInsight Manager系统,选择“集群 > 服务 > Kafka > 配置 > 全部配置”,搜索“allow.everyone.if.no.acl.found”配置,修改参数值为true,保存配置即可。
- 示例2:为左表和右表分别设置过期时间,右表可以是子查询。
INSERT INTO print SELECT t1.user_id, t1.user_name, t3.score FROM t1 JOIN -- 为左表和右表设置不同的TTL时间 /*+ OPTIONS('state.ttl.left' = '60S', 'state.ttl.right' = '120S') */ ( select UPPER(t2.user_id) as user_id, t2.score from t2 ) as t3 ON t1.user_id = t3.user_id;
常见问题
Q:state.ttl.left和state.ttl.right是否可以只配置其中一个?
A:可以只配置其中一个,但未配置的一侧将使用表级别的全局TTL配置(即table.exec.state.ttl参数的值)。建议根据业务需求同时配置两侧的TTL,以获得最佳的状态清理效果。