Flink SQL
Flink SQL支持动态表上复杂而灵活的连接操作。有几种不同类型的Join来满足可能需要的各种语义查询。
默认情况下,连接的顺序没有优化。表按照在FROM子句中指定的顺序连接。您可以调整join方式提高查询的性能,方法是先列出更新频率最低的表,最后列出更新频率最高的表。请确保不产生交叉连接(笛卡尔积),交叉连接不受支持,会导致查询失败。
1、Regular Joins
常规连接是最通用的连接类型,其中任何新记录或连接两端的更改都是可见的,并且会影响整个连接结果。例如,如果左侧有一个新记录,当产品id相等时,它将与右侧所有以前和将来的记录连接。
Regular Join会将两个关联表长存再状态中,会导致状态越来越大。
-- 创建学生表流表,数据再kafka中
CREATE TABLE student_join (
id String,
name String,
age int,
gender STRING,
clazz STRING
) WITH (
'connector' = 'kafka',
'topic' = 'student_join',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'csv',
'scan.startup.mode' = 'latest-offset'
);
-- 分数表
CREATE TABLE score_join (
s_id String,
c_id String,
sco int
) WITH (
'connector' = 'kafka',
'topic' = 'score_join',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'csv',
'scan.startup.mode' = 'latest-offset'
);
--- inner join
select a.id,a.name,b.sco from
student_join as a
inner join
score_join as b
on a.id=b.s_id
-- left outer join
select a.id,a.name,b.sco from
student_join as a
left join
score_join as b
on a.id=b.s_id
-- full outer join
select a.id,a.name,b.sco from
student_join as a
full join
score_join as b
on a.id=b.s_id
-- 创建生产者向两个topic中生产数据
kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic score_join
1500100001,1000001,98
1500100001,1000002,5
1500100001,1000003,0
kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic student_join
1500100001,施笑槐,22,女,文科六班
1500100002,吕金鹏,24,男,文科七班
2、Interval Joins
返回一个受连接条件和时间限制限制的简单笛卡尔积。Interval Joins至少需要一个等值连接条件和一个限制两边时间的连接条件。
-- 创建学生表流表,数据再kafka中
CREATE TABLE student_join_proc (
id String,
name String,
age int,
gender STRING,
clazz STRING,
stu_time as PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'student_join',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'csv',
'scan.startup.mode' = 'latest-offset'
);
-- 分数表
CREATE TABLE score_join_proc (
s_id String,
c_id String,
sco int,
sco_time as PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'score_join',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'csv',
'scan.startup.mode' = 'latest-offset'
);
-- Interval Joins
select a.id,a.name,b.sco from
student_join_proc as a, score_join_proc as b
where a.id=b.s_id
and a.stu_time BETWEEN b.sco_time - INTERVAL '10' SECOND AND b.sco_time
-- 创建生产者向两个topic中生产数据
kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic score_join
1500100001,1000001,98
1500100001,1000002,5
1500100002,1000003,0
kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic student_join
1500100001,施笑槐,22,女,文科六班
1500100002,吕金鹏,24,男,文科七班
3、Temporal Joins
时态表是一个随时间变化的表,在Flink中也称为动态表。时态表中的行与一个或多个时态周期相关联,所有Flink表都是时态表(动态的)。时态表包含一个或多个版本表快照,它可以是一个跟踪更改的历史表。数据库变更日志,包含所有快照或一个具体化变更的变化维度表(例如。包含最新快照的数据库表)。
-- 订单表
CREATE TABLE orders (
order_id STRING, -- 订单编号
price DECIMAL(32,2), --订单的价格
currency STRING, -- 汇率表主键
order_time TIMESTAMP(3), -- 订单发生的事件
WATERMARK FOR order_time AS order_time -- 设置事件时间和水位线
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'csv',
'scan.startup.mode' = 'latest-offset'
);
--汇率表
CREATE TABLE currency_rates (
currency STRING, -- 汇率表主键
conversion_rate DECIMAL(32, 2), -- 汇率
update_time TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL, --汇率更新时间
WATERMARK FOR update_time AS update_time,--时间字段和水位线
PRIMARY KEY(currency) NOT ENFORCED--设置主键
) WITH (
'connector' = 'kafka',
'topic' = 'bigdata.currency_rates',
'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092',
'properties.group.id' = 'asdasdasd',
'format' = 'canal-json',
'scan.startup.mode' = 'earliest-offset',
'canal-json.ignore-parse-errors' = 'true'
);
-- Temporal Joins
SELECT
order_id,
price,
orders.currency,
conversion_rate,
order_time
FROM orders
LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF orders.order_time
ON orders.currency = currency_rates.currency;
-- 订单表数据
kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic orders
001,1000.0,1001,2022-08-02 11:08:20
001,1000.0,1001,2022-08-02 11:15:55
001,1000.0,1001,2022-08-02 11:20:55 4、流表(kafka)关联维表(hbase,mysql)
1、Regular Join
使用常规jon做维表关联,会出现数据库中维表更新了,但是flink中无法捕获更新,只能关联到任务刚启动时读取的数据
2、Lookup Join
Lookup Join是一种查找join,当流表中摄入一行数据时会使用关联字段倒维表的数据源中查询数据,这个方式可以保证每次都能关联到最新的维度信息,当然性能会有一点的降低,可以通过增加缓存来优化性能
Flink SQL
Flink SQL支持动态表上复杂而灵活的连接操作。有几种不同类型的Join来满足可能需要的各种语义查询。
默认情况下,连接的顺序没有优化。表按照在FROM子句中指定的顺序连接。您可以调整join方式提高查询的性能,方法是先列出更新频率最低的表,最后列出更新频率最高的表。请确保不产生交叉连接(笛卡尔积),交叉连接不受支持,会导致查询失败。
1、Regular Joins
常规连接是最通用的连接类型,其中任何新记录或连接两端的更改都是可见的,并且会影响整个连接结果。例如,如果左侧有一个新记录,当产品id相等时,它将与右侧所有以前和将来的记录连接。
Regular Join会将两个关联表长存再状态中,会导致状态越来越大。
2、Interval Joins
返回一个受连接条件和时间限制限制的简单笛卡尔积。Interval Joins至少需要一个等值连接条件和一个限制两边时间的连接条件。
3、Temporal Joins
时态表是一个随时间变化的表,在Flink中也称为动态表。时态表中的行与一个或多个时态周期相关联,所有Flink表都是时态表(动态的)。时态表包含一个或多个版本表快照,它可以是一个跟踪更改的历史表。数据库变更日志,包含所有快照或一个具体化变更的变化维度表(例如。包含最新快照的数据库表)。
-- 订单表 CREATE TABLE orders ( order_id STRING, -- 订单编号 price DECIMAL(32,2), --订单的价格 currency STRING, -- 汇率表主键 order_time TIMESTAMP(3), -- 订单发生的事件 WATERMARK FOR order_time AS order_time -- 设置事件时间和水位线 ) WITH ( 'connector' = 'kafka', 'topic' = 'orders', 'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092', 'properties.group.id' = 'asdasdasd', 'format' = 'csv', 'scan.startup.mode' = 'latest-offset' ); --汇率表 CREATE TABLE currency_rates ( currency STRING, -- 汇率表主键 conversion_rate DECIMAL(32, 2), -- 汇率 update_time TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL, --汇率更新时间 WATERMARK FOR update_time AS update_time,--时间字段和水位线 PRIMARY KEY(currency) NOT ENFORCED--设置主键 ) WITH ( 'connector' = 'kafka', 'topic' = 'bigdata.currency_rates', 'properties.bootstrap.servers' = 'master:9092,node1:9092,node2:9092', 'properties.group.id' = 'asdasdasd', 'format' = 'canal-json', 'scan.startup.mode' = 'earliest-offset', 'canal-json.ignore-parse-errors' = 'true' ); -- Temporal Joins SELECT order_id, price, orders.currency, conversion_rate, order_time FROM orders LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF orders.order_time ON orders.currency = currency_rates.currency; -- 订单表数据 kafka-console-producer.sh --broker-list master:9092,node1:9092,node2:9092 --topic orders 001,1000.0,1001,2022-08-02 11:08:20 001,1000.0,1001,2022-08-02 11:15:55 001,1000.0,1001,2022-08-02 11:20:554、流表(kafka)关联维表(hbase,mysql)
1、Regular Join
使用常规jon做维表关联,会出现数据库中维表更新了,但是flink中无法捕获更新,只能关联到任务刚启动时读取的数据
2、Lookup Join
Lookup Join是一种查找join,当流表中摄入一行数据时会使用关联字段倒维表的数据源中查询数据,这个方式可以保证每次都能关联到最新的维度信息,当然性能会有一点的降低,可以通过增加缓存来优化性能