0%

Flink-SQL实践

Flink-SQL-active

维表同步脚本:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert
into
dim_result_lct_activy_config
select
Fact_id,
LAST_VALUE(Fact_name),
FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_start_time)),
'yyyyMMddHHmmss') as startTime,
FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_end_time)),
'yyyyMMddHHmmss') as endTime,
LAST_VALUE(Fstate)
from
db_act_config_t_act_logic_config
group by
Fact_id

FlinkSQL Upsert/Retraction 写入 MySQL 的问题

retract流是算子的特性,不同的算子不同,比如group by聚合算子天生就是会发retract流的,比如你之前算出来的是a, 1, 后面变化了,变成了a,2,所以group by算子会发-a,1, +a,2.

像join的话也是有retract流的,比如outer join

当然如果下游sink支持upsert, 这块也会有优化,会优化为update。

INSERT INTO mysql_sink SELECT f1, count(*) FROM kafka_src GROUP BY f1
每从 kafka 过来一条新的记录,会生成两条记录 Tuple2<Row, Boolean>, 旧的被删除,新的会添加上。这是query是会一个会产生retract stream的query,可以简单理解成每条kafka的数据过来会产生两条记录,但是最终写入下游的系统。需要看下游的系统支持和实现的sink(现在有三种sink AppendStreamSink, UpsertStreamSink, RetractStreamSink)

我看 https://github.com/apache/flink/tree/master/flink-connectors/flink-jdbc/src/main/java/org/apache/flink/api/java/io/jdbc 没有 Retract 方式
实际上使用了 JDBCUpsertTableSink.java 的代码写入 MySQL 吗?
现有的sink中,kafka是实现的AppendStreamSink,所以只支持insert 的记录,不支持retract.
你用DDL声明的mysql表,对应的jdbc sink 是JDBCUpsertTableSink,所以会按照Upsert的逻辑处理, 也不支持retract。

如若不带 group by 直接:
INSERT INTO mysql_sink SELECT f1, f2 FROM kafka_src
主键冲突写入 mysql 是会出错的,怎么可以用 Upsert 的方式直接覆盖呢?

不带 group by时无法推导出query的 unique key,没法做按照unique key的更新,
只需要将 query的 key (你这里是group by 后的字段)和db中主键保持一致即可

TopN

1
2
3
4
5
6
7
SELECT [column_list]
FROM (
SELECT [column_list],
ROW_NUMBER() OVER ([PARTITION BY col1[, col2...]]
ORDER BY col1 [asc|desc][, col2 [asc|desc]...]) AS rownum
FROM table_name)
WHERE rownum <= N [AND conditions]
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
EXPLAIN PLAN FOR SELECT * FROM (
SELECT *,
row_number() OVER(PARTITION BY merchandiseId ORDER BY totalQuantity DESC) AS rownum
FROM (
SELECT merchandiseId, sum(quantity) AS totalQuantity
FROM rtdw_dwd.kafka_order_done_log
GROUP BY merchandiseId
)
) WHERE rownum <= 10

== Abstract Syntax Tree ==
LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[$2])
+- LogicalFilter(condition=[<=($2, 10)])
+- LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[ROW_NUMBER() OVER (PARTITION BY $0 ORDER BY $1 DESC NULLS LAST)])
+- ...

== Optimized Logical Plan ==
Rank(strategy=[RetractStrategy], rankType=[ROW_NUMBER], rankRange=[rankStart=1, rankEnd=10], partitionBy=[merchandiseId], orderBy=[totalQuantity DESC], select=[merchandiseId, totalQuantity, w0$o0])
+- Exchange(distribution=[hash[merchandiseId]])
+- ...

== Physical Execution Plan ==
Stage 1 : Data Source
...

Stage 2 : Operator
...

Stage 4 : Operator
...

Stage 6 : Operator
content : Rank(strategy=[RetractStrategy], rankType=[ROW_NUMBER], rankRange=[rankStart=1, rankEnd=10], partitionBy=[merchandiseId], orderBy=[totalQuantity DESC], select=[merchandiseId, totalQuantity, w0$o0])
ship_strategy : HASH

滚动窗口 TUMBLE window

每分钟pv、uv

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
-- 滚动窗口统计pv:计算每分钟用户的页面访问数量
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
TUMBLE_START(up.fevent_time, INTERVAL '1' MINUTE) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv,
count(distinct fuin) as uv
FROM
datagen AS up
GROUP BY
-- 按照1分钟进行开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
TUMBLE(up.fevent_time, INTERVAL '1' MINUTE) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据
1
2
3
4
5
6
7
8
9
10
11
12
INSERT INTO
t_hj_sample_result_040901
SELECT
CAST(up.Fact_id AS INTEGER) AS catalog_01,
CAST(
TUMBLE_END(up.ptime, INTERVAL '10' MINUTE) AS BIGINT
) AS window_end
FROM
b_cdg_cft_lct_act_data_db_act_t_user_prize AS up
GROUP BY
CAST(up.Fact_id AS INTEGER),
TUMBLE(up.ptime, INTERVAL '10' MINUTE)

滑动窗口

每小时统计一次过去一天的pv

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
-- 滑动窗口统计pv:每小时统计一次过去1天的用户页面访问数量
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
HOP_START(
up.fevent_time,
INTERVAL '1' HOUR,
INTERVAL '1' DAY
) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv
FROM
datagen AS up
GROUP BY
-- 按照1小时进行滑动开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定统计多久之前的数据
HOP(
up.fevent_time,
INTERVAL '1' HOUR,
INTERVAL '1' DAY
) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据

每小时统计一次过去一周的各页面的pv

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
-- 滑动窗口统计uv:每天统计一次过去一周网页的独立用户访问量
-- 上游表为tube,写入mysql
-- 假设上游表为用户交易表,计算只需要对用户id进行去重即可
INSERT INTO
console
SELECT
-- 窗口的开始时间
FROM_UNIXTIME(
CAST(
HOP_START(up.fevent_time, INTERVAL '1' DAY, INTERVAL '7' DAY) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
-- 使用distinct算子进行去重
count(distinct(fuin)) as uv
From
datagen AS up
GROUP BY
-- 按照一天进行开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定统计多久之前的数据
HOP(up.fevent_time, INTERVAL '1' DAY, INTERVAL '7' DAY) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据

增量窗口 INCREMENT

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
-- 增量窗口统计pv:每天统计一次用户页面访问数量,同时每分钟输出一次实时统计结果
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
INCREMENT_START(
up.fevent_time,
INTERVAL '1' DAY,
INTERVAL '1' HOUR
) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv
FROM
datagen AS up
GROUP BY
-- 按照1小时进行滑动开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定多久输出一次实时结果
INCREMENT(
up.fevent_time,
INTERVAL '1' DAY,
INTERVAL '1' HOUR
) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据