0%

Flink SQL中的Join

Flink SQL Join

Join的几种形式,与实现原理

双流Join

Left join

TTL

从 Flink 1.6 版本开始,社区引入了状态 TTL(Time-To-Live)特性。在通过Flink SQL 实现流处理时,开发者可以为作业 SQL 设置TTL,实现过期状态的自动清理,从而防止作业状态无限膨胀

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
SELECT
t_date,
COUNT ( DISTINCT user_id ) AS cnt_login, -- 今日登录用户数
COUNT ( DISTINCT CASE WHEN t_date = t_debut THEN user_id END ) AS cnt_new -- 今日新用户数
FROM
(
-- 计算每个用户有史以来的最小登录时间
SELECT
t_date,
user_id,
MIN (t_date) OVER (
PARTITION BY user_id
ORDER BY proctime
ROWS BETWEEN 1 PRECEDING AND CURRENT ROW
) AS t_debut
FROM Login
) AS t
GROUP BY t_date

Query Configuration 查询配置

https://ci.apache.org/projects/flink/flink-docs-stable/dev/table/streaming/query_configuration.html

Flink Table API 和SQL接口提供参数来调整连续查询的准确性和资源消耗。参数通过 QueryConfig 对象指定。QueryConfig 可以从 TableEnvironment 获得。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 获取 query configuration
val qConfig: StreamQueryConfig = tableEnv.queryConfig

// 设置查询参数
qConfig.withIdleStateRetentionTime(Time.hours(12), Time.hours(24))

// 定义查询和 TableSink
val result: Table = ???
val sink: TableSink[Row] = ???

// TableSink 发送结果表时传递查询参数
result.writeToSink(sink, qConfig)

// 转换为 DataStream 时传递查询参数
val stream: DataStream[Row] = result.toAppendStream[Row](qConfig)

空闲状态保持时间(Idle State Retention Time)参数定义一个键的状态在一次更新之后保存多久后删除。

通过删除键的状态,连续查询会完全忘记它之前已经看过这个键。如果删除的键再次出现,则被视为具有相应键的第一个记录。对于前面的查询示例,这意味着 sessionId 的计数从0开始。

配置空闲状态保存时间有两个参数:

minimum idle state retention time,定义非活动键的状态在删除前至少保持多少时间。
maximum idle state retention time,定义非活动键的状态在删除前最多保持多少时间。

对于前面的查询示例:

1
2
3
4
val qConfig: StreamQueryConfig = ???

// 设置 idle state retention time: min = 12 hours, max = 24 hours
qConfig.withIdleStateRetentionTime(Time.hours(12), Time.hours(24))

清理状态需要额外的记录,对于 minTime 和 maxTime 较大差异的情况成本更低,因此 minTime 和 maxTime 直接必须至少相差5分钟。

维表Join

启用AsyncIO

时间表Join: Temporal Table Join

1
2
3
4
5
6
SELECT
o.amout, o.currency, r.rate, o.amount * r.rate
FROM
Orders AS o
JOIN LatestRates FOR SYSTEM_TIME AS OF o.proctime AS r
ON r.currency = o.currency