flink中sql关键字一不小心使用了状态,你真的都知道吗
·
0. 先给总结:
以下操作会引入状态:
(1)聚合:SUM, COUNT, AVG, MAX, MIN 等。
(2)窗口:TUMBLE, HOP, SESSION, CUMULATE 等。
(3)JOIN:时间窗口 Join、流对表 Join。
(4)DISTINCT去重操作 和 GROUP BY 分组操作。
(5)OVER 窗口:滑动窗口聚合。
(6)TOP N:ROW_NUMBER(), RANK() 等。
(7)自定义 UDF/UDAF。
(8)使用Lookup Join 开启缓存
============
状态引入的优化建议
(1)减少状态量:通过清理过期数据、优化窗口大小、减少分组键等方式控制状态大小。
(2)使用 RocksDB State Backend:对于大状态量,使用持久化的状态后端。
(3)监控状态大小:定期检查 Flink 的状态指标,防止状态过大导致作业崩溃。
1. 聚合操作 (Aggregations)
聚合操作需要维护中间结果,因此会引入状态。常见的聚合操作包括:
- SUM: 累加求和,需要保存当前累加值。
- COUNT: 计数操作,需要保存当前的计数值。
- AVG: 平均值,需要保存累加和以及计数。
- MAX/MIN: 最大值或最小值,需要保存当前的最大/最小值。
示例
SELECT
key,
SUM(value) AS total
FROM table_name
GROUP BY key;
2. 窗口操作 (Windows)
窗口操作需要在窗口范围内保存事件数据或中间状态,以支持流式计算。
窗口类型:
- TUMBLE: 滚动窗口(固定大小,不重叠)。
- HOP: 滑动窗口(可以有重叠,窗口开始时间间隔小于窗口长度)。
- SESSION: 会话窗口(根据活动间隔自动调整窗口范围)。
- CUMULATE: 累积窗口(窗口逐步扩大,直到包含所有数据)。
示例
SELECT
TUMBLE_START(rowtime, INTERVAL '1' MINUTE) AS window_start,
COUNT(*) AS cnt
FROM table_name
GROUP BY TUMBLE(rowtime, INTERVAL '1' MINUTE);
3. JOIN 操作
Join 操作会引入状态,因为它需要将两张表的数据缓存在状态中以完成匹配:
- 流对表(Temporal Table)Join。
- 流对流(Stream-to-Stream)Join。
示例
SELECT
stream1.key,
stream2.value
FROM stream1
JOIN stream2
ON stream1.key = stream2.key
AND stream1.rowtime BETWEEN stream2.rowtime - INTERVAL '5' SECOND AND stream2.rowtime + INTERVAL '5' SECOND;
4. OVER 窗口 (Over Windows)
OVER 窗口需要维护窗口范围内的数据状态,用于计算滑动窗口内的累积结果。
示例
SELECT
key,
SUM(value) OVER (PARTITION BY key ORDER BY rowtime ROWS BETWEEN 5 PRECEDING AND CURRENT ROW) AS sliding_sum
FROM table_name;
4. OVER 窗口 (Over Windows)
OVER 窗口需要维护窗口范围内的数据状态,用于计算滑动窗口内的累积结果。
示例
SELECT
key,
SUM(value) OVER (PARTITION BY key ORDER BY rowtime ROWS BETWEEN 5 PRECEDING AND CURRENT ROW) AS sliding_sum
FROM table_name;
5. DISTINCT 和GROUP BY操作
distinct 操作需要保存已经看到的所有值的集合,以防止重复,因此会引入状态。
group by 操作需要为每个分组维护状态信息,包括每个分组的聚合中间值。
示例
SELECT
DISTINCT key
FROM table_name;
===============
SELECT
key,
COUNT(*) AS cnt
FROM table_name
GROUP BY key;
6. TOP N (Rank/ROW_NUMBER 等)
排序和 Top N 查询需要缓存和排序大量数据,因此会引入状态。
- ROW_NUMBER(): 需要维护分区内的排序状态。
- RANK() 和 DENSE_RANK(): 同样需要维护排序状态。
示例
SELECT
key,
ROW_NUMBER() OVER (PARTITION BY category ORDER BY score DESC) AS rank
FROM table_name;
7. 自定义 UDF 和 UDAF
自定义函数可能会引入状态,尤其是用户在函数中使用了 Flink 的状态管理 API,例如:
- ValueState
- ListState
- MapState
示例
CREATE FUNCTION my_aggregate_function AS 'com.example.MyUDAF';
SELECT
key,
my_aggregate_function(value)
FROM table_name
GROUP BY key;
8.Lookup Join
在 Flink 中,Lookup Join 通常用于将流数据与维表数据(例如数据库、缓存等外部存储系统)进行关联。为了高效处理查询,Flink 会对维表数据进行缓存(State),从而避免每次都需要从外部系统中查询维表数据。这种缓存机制会引入状态。
-- 使用 Lookup Join 查询订单与客户信息
SELECT
o.order_id,
o.customer_id,
c.customer_name,
o.order_time
FROM orders AS o
LEFT JOIN dim_customer FOR SYSTEM_TIME AS OF o.order_time AS c
ON o.customer_id = c.customer_id;
更多推荐



所有评论(0)