0%

Flink-语法糖

Flink-SQL语法

Apache Flink SQL training

group window

groupByWindow会直接生成回撤流

1
2
3
4
5
6
7
8
9
10
11
12
13
insert       
into
dim_result_lct_activy_config
select
Fact_id,
LAST_VALUE(Fact_name),
regexp_Replace( LAST_VALUE(Fact_start_time), '-|:|\s','') as startTime,
regexp_Replace(LAST_VALUE(Fact_end_time),'-|:|\s','') as endTime,
LAST_VALUE(Fstate)
from
db_act_config_t_act_logic_config
group by
Fact_id

这是一个同步数据的demo,db_act_config_t_act_logic_config 是kafka数据源,来自源MySQL的变更数据;dim_result_lct_activy_config是目的表,Fact_id为主键。

group window生成retract stream

1
insert into mysql_sink select fkey,count(1) as cnt from kafka_source group by fkey

上述语句是一个group window, 每从kafka中过来一条数据,都会产生两条记录(Tuple2<Row,Boolean>), 删除旧记录,添加新记录。

group window会产生 retract stream, 下游系统必须支持retract stream,(目前共有三种sink: AppendStreamSink, UpsertStreamSink, RetractStreamSink )

Flink-connector-JDBC 使用的是JDBCUpsertTableSink.java写入MySQL, 支持Retract

https://github.com/apache/flink/tree/master/flink-connectors/flink-jdbc/src/main/java/org/apache/flink/api/java/io/jdbc

Flink-connector-kafka 实现的是 AppendStreamSink,只支持insert,不支持retract. 会报错

AppendStreamTableSink requires that Table has only insert changes

1
insert into mysql_sink select fkey,count(1) as cnt from kafka_source

如果不带group by, 无法推导出unique key, 无法按照unique key更新

http://apache-flink.147419.n8.nabble.com/FlinkSQL-Upsert-Retraction-MySQL-td2785.html

1
2
3
4
5
6
7
8
9
10
11
/** 
* Get dialect upsert statement, the database has its own upsert syntax, such as Mysql
* using DUPLICATE KEY UPDATE, and PostgresSQL using ON CONFLICT... DO UPDATE SET..
*
* @return None if dialect does not support upsert statement, the writer will degrade to
* the use of select + update/insert, this performance is poor.
*/
default Optional<String> getUpsertStatement(
String tableName, String[] fieldNames, String[] uniqueKeyFields) {
return Optional.empty();
}

不同的数据库产品有不同的语句,所以默认实现是delete +insert

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Override 
public void executeBatch() throws SQLException {
if (keyToRows.size() > 0) {
for (Map.Entry<Row, Tuple2<Boolean, Row>> entry : keyToRows.entrySet()) {
Row pk = entry.getKey();
Tuple2<Boolean, Row> tuple = entry.getValue();
if (tuple.f0) {
processOneRowInBatch(pk, tuple.f1);
} else {
setRecordToStatement(deleteStatement, pkTypes, pk);
deleteStatement.addBatch();
}
}
internalExecuteBatch();
deleteStatement.executeBatch();
keyToRows.clear();
}
}

image-20210928204455343

image-20210928204646114

image-20210928204703778

Over window

SQL窗口函数 传统SQL窗口函数的介绍

1
2
3
4
5
6
7
8
9
10
11
12
13
14
select 
to_char(SYSTIMESTAMP(),'yyyymmddhh24miss') fetl_time,
*
from
(
select
*,
row_number() over(partition by fid order by fmodify_time desc,exp_time_stample_order desc) rn
from
db.table1
where
fdate=20210101
) t
where rn=1
1
2
3
4
5
6
7
8
9
10
11
SELECT COUNT(amount) OVER (
PARTITION BY user
ORDER BY proctime
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW)
FROM Orders;
SELECT COUNT(amount) OVER w, SUM(amount) OVER w
FROM Orders
WINDOW w AS (
PARTITION BY user
ORDER BY proctime
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) ;

OVER 窗口应用示例

首先通过 DDL 定义源数据表和结果表,如下输入是用户行为消息,输出到计算结果消息。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
CREATE TABLE `user_action` (
`user_id` VARCHAR,
`page_id` VARCHAR,
`action_type` VARCHAR,
`event_time` TIMESTAMP,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector.type' = 'kafka',
'connector.topic' = 'user_action',
'connector.version' = '0.11',
'connector.properties.0.key' = 'bootstrap.servers',
'connector.properties.0.value' = 'xxx:9092',
'connector.startup-mode' = 'latest-offset',
'update-mode' = 'append',
'...' = '...'
);

CREATE TABLE `agg_result` (
`user_id` VARCHAR,
`page_id` VARCHAR,
`result_type` VARCHAR,
`result_value` BIGINT
) WITH (
'connector.type' = 'kafka',
'connector.topic' = 'agg_result',
'...' = '...'
);

场景一,实时触发的最近2小时用户+页面维度的点击量,注意窗口是向前2小时,类似于实时触发的滑动窗口。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert into
agg_result
select
user_id,
page_id,
'click-type1' as result_type
count(1) OVER (
PARTITION BY user_id, page_id
ORDER BY event_time
RANGE BETWEEN INTERVAL '2' HOUR PRECEDING AND CURRENT ROW
) as result_value
from
user_action
where
action_type = 'click'

场景二,实时触发的当天用户+页面维度的浏览量,这就是开篇问题解法,其中多了一个日期维度分组条件,这样就做到输出结果从滑动时间转为固定时间(根据时间区间分组),因为 WATERMARK 机制,今天并不会有昨天数据到来(如果有都被自动抛弃),因此只会输出今天的分组结果。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert into
agg_result
select
user_id,
page_id,
'view-type1' as result_type
count(1) OVER (
PARTITION BY user_id, page_id, DATE_FORMAT(event_time, 'yyyyMMdd')
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW
) as result_value
from
user_action
where
action_type = 'view'

场景三,实时触发的当天用户+页面点击率 CTR(Click-Through-Rate),这相比前面增加了多个 OVER 聚合计算,可以将窗口定义写在最后。注意示例中缺少了类型转换,因为除法结果是 decimal,也缺少精度处理函数 ROUND。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
insert into
agg_result
select
user_id,
page_id,
'ctr-type1' as result_type,
sum(
case
when action_type = 'click' then 1 else 0
end
) OVER w
/
if(
sum(
case
when action_type = 'view' then 1 else 0
end
) OVER w = 0,
1,
sum(
case
when action_type = 'view' then 1 else 0
end
) OVER w
)
as result_value
from
user_action
where
1 = 1
WINDOW w AS (
PARTITION BY user_id, page_id, DATE_FORMAT(event_time,'yyyyMMdd')
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW
)

此外,OVER 窗口聚合还可以支持查询子句、关联查询、UNION ALL 等组合,并可以实现对关联出来的列进行聚合等复杂情况。

实时TopN

SQL实时TopN

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
create table source_kafka 
(
userID String,
eventType String,
eventTime String,
productID String
) with (
'connector.type' = 'kafka',
'connector.version' = '0.10',
'connector.properties.bootstrap.servers' = 'kafka01:9092',
'connector.properties.zookeeper.connect' = 'kafka01:2181',
'connector.topic' = 'test_1',
'connector.properties.group.id' = 'c1_test_1',
'connector.startup-mode' = 'latest-offset',
'format.type' = 'json'
);
create table sink_mysql
(
datetime STRING,
productID STRING,
userID STRING,
clickPV BIGINT
) with (
'connector.type' = 'jdbc',
'connector.url' = 'jdbc:mysql://localhost:3306/bigdata',
'connector.table' = 't_product_click_topn',
'connector.username' = 'root',
'connector.password' = 'bigdata',
'connector.write.flush.max-rows' = '50',
'connector.write.flush.interval' = '2s',
'connector.write.max-retries' = '3'
);
INSERT INTO sink_mysql
SELECT datetime, productID, userID, clickPV
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY datetime, productID ORDER BY clickPV desc) AS rownum
FROM (
SELECT SUBSTRING(eventTime,1,13) AS datetime,
productID,
userID,
count(1) AS clickPV
FROM source_kafka
GROUP BY SUBSTRING(eventTime,1,13), productID, userID
) a
) t
WHERE rownum <= 3;

OVER 窗口问题和优化

在底层实现中,所有细分 OVER 窗口的数据都是共享的,只存一份,这点不像滑动窗口会保存多份窗口数据。但是 OVER 窗口会把所有数据明细存在状态后端中(内存、RocksDB 或 HDFS),每一次窗口计算后会清除过期数据。因此如果向前窗口时间较大,或数据明细过多,可能会占用大量内存,即使通过 RocksDB 存在磁盘上,也有因为磁盘访问慢导致性能下降进而产生反压问题。在实现源码 RowTimeRangeBoundedPrecedingFunction 可以看到虽然每次窗口计算时新增聚合值和减少过期聚合值是增量式的,不用遍历全部窗口明细,但是为了计算过期数据,即超过 PRECEDING 的数据,仍然需要把存储的那些时间戳全部拿出来遍历,判断是否过期,以及是否要减少聚合值。我们尝试了通过数据有序性减少查询操作,但是效果并不明显,目前主要是配置调优和加大任务分片数进行优化。