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;

更多推荐