Flink sql —— Lookup Join
一、背景
实时生产中,常常有关联维表的需求,如关联资源维表获取资源信息。但维表不是一成不变的,会不断变化。如果用常规的 JOIN 语法来表达关联维表,是无法知晓关联上的是哪个时刻的维表呢?
二、方案
Flink SQL 的维表 JOIN 语法引入了Temporal Table 的标准语法,用来声明关联的是维表哪个时刻的快照。
官网地址:https://nightlies.apache.org/flink/flink-docs-release-2.1/docs/dev/table/sql/queries/joins/
SELECT t.app_id,
a.app_name
FROM event_inc_inc_rt AS t
LEFT JOIN dim_app FOR SYSTEM_TIME AS OF t.proc_time AS a
ON t.app_id = a.app_id
如上语法所示的,维表后面要跟上 FOR SYSTEM_TIME AS OF PROCTIME() 关键字,其含义是每条到达的流式数据所关联上的是到达时刻的维表快照。这里的时间可以是事件时间/处理时间。

实现机制:通过周期性或按需查询外部存储(如 MySQL、HBase)获取维表数据。当主流数据到达时,通过 Key 查询本地缓存或直接访问外部数据库,将维表属性补充到流数据中。
注意:主表后面进来的数据只会关联当时维表的最新信息,如果 JOIN 行为发生后,维表中的数据发生了变化(新增、更新或删除),则已关联的维表数据不会被同步变化。
三、调优
维表 Join 的调优方向:
① 提高吞吐量,维表的IO 请求有可能会阻塞实时数据流的分析计算;—— 异步
② 降低维表数据库读取次数,否则频繁查询会造成一定程度的延迟;—— 缓存
某数据平台redis表配置:

异步lookup数据库关联

相较于同步lookup,异步方式可大大提高数据库查询的吞吐量,但相应的也会加大数据库的负载,并且由于查询只能查当前时间点的维度数据,因此可能造成数据查询结果的不准确。
缓存

当每个数据进来时,先去本地缓存中查询,如果存在则直接关联输出,减少了一次 IO 请求。如果不存在,再发起数据库查询请求,请求返回的结果会先存入缓存中以备下次查询。
问题:缓存的数据无法及时更新,可能会造成关联数据不正确。
缓存策略
Flink SQL 目前提供两种缓存策略,LRU cache 和 ALL cache。
| 缓存策略 | 机制 | 适用场景 |
|---|---|---|
| ALL cache | 将整个维表缓存到本地 | 小维表 join |
| LRU cache(Least Recently Used) | 淘汰最近最少被访问的数据 | 超大维表 join |
缓存参数
| 参数 | 含义 | 作用 | 生效机制 | 生效策略 |
|---|---|---|---|---|
| cacheTTLMs | 缓存失效时间 | 定期更新维表数据 | 作用于每条缓存数据上的,也就是某条缓存数据在指定 timeout 时间内没有被访问,则会从缓存中移除 | ALL cache/LRU cache |
| cacheSize | 缓存的最大数据行数 | 超过cacheSize 行数据时,会根据 LRU算法进行淘汰,以保证缓存中数据的有效性和时效性 | LRU cache |
cacheTTLMs 用于标识数据是否“失效”,但不会主动删除未访问的过期数据。 只有满足以下条件之一时,数据 A 才会被真正移除:
① 缓存条目总数超过 max-rows,且数据 A 是最近最少使用的(LRU 淘汰) ;
② 数据 A 被再次访问时发现已过期,此时会被剔除并重新查询 Redis;
如果数据 A 既未达到访问上限,又未被再次访问,即使它已经过期,也仍然会保留在缓存中(但查询时会自动发现过期并更新)
更多推荐
所有评论(0)