Fluss初窥四与Flinksql集成使用
Flink DDL
创建 Fluss Catalog
sql
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'fluss-server-1:9123'
);
USE CATALOG fluss_catalog;
以下属性在使用 Fluss 目录时可以设置:

创建数据库
sql
CREATE DATABASE my_db;
USE my_db;
删除数据库
sql
-- Flink 不允许删除当前数据库,切换到 Fluss 默认数据库
USE fluss;
-- 删除数据库
DROP DATABASE my_db;
创建主键表
sql
CREATE TABLE my_pk_table (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (shop_id, user_id) NOT ENFORCED
) WITH (
'bucket.num' = '4'
);
创建日志表
sql
CREATE TABLE my_log_table (
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING
) WITH (
'bucket.num' = '8'
);
创建分区主键表
分区表必须启用自动分区并设置自动分区时间单位。
sql
CREATE TABLE my_part_pk_table (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (dt, shop_id, user_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket.num' = '4',
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'day'
);
创建分区日志表
sql
CREATE TABLE my_part_log_table (
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING,
dt STRING
) PARTITIONED BY (dt) WITH (
'bucket.num' = '8',
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'hour'
);
建表时,“with”参数支持的选项如下:

表配置参数:

客户端配置参数:

使用 LIKE 语法复制表
创建具有与另一个表相同模式、分区和表属性的表的语句,请使用 CREATE TABLE LIKE。
sql
-- 创建一个临时 datagen 表
CREATE TEMPORARY TABLE datagen (
user_id BIGINT,
item_id BIGINT,
behavior STRING,
dt STRING,
hh STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '10'
);
-- 创建一个 Fluss 表,从临时表派生元数据,排除选项
CREATE TABLE my_table LIKE datagen (EXCLUDING OPTIONS);
删除表
sql
DROP TABLE my_table;
Flink Writer
可以直接使用 INSERT INTO 语句将数据插入或更新到 Fluss 表中。Fluss 主键表可以接受所有类型的消息(INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE),而 Fluss 日志表只能接受 INSERT 类型消息。
插入
INSERT INTO 语法用于将数据写入 Fluss 表。它支持流式和批量模式,并与主键表(用于更新数据)以及日志表(用于追加数据)兼容。
追加数据到日志表
sql
-- 创建日志表(target)
CREATE TABLE log_table (
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING
);
-- 生成临时表
CREATE TEMPORARY TABLE source (
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING
) WITH ('connector' = 'datagen');
-- 插入数据
INSERT INTO log_table
SELECT * FROM source;
主键表的更新操作
sql
-- 创建主键表(target)
CREATE TABLE pk_table (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (shop_id, user_id) NOT ENFORCED
);
-- 生成临时表
CREATE TEMPORARY TABLE source (
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT
) WITH ('connector' = 'datagen');
-- 更新所有列
INSERT INTO pk_table
SELECT * FROM source;
-- 更新部分列
INSERT INTO pk_table (shop_id, user_id, num_orders)
SELECT shop_id, user_id, num_orders FROM source;
删除
Fluss 支持通过 DELETE FROM 语句在批处理模式下删除主键表的数据。目前,只支持基于主键的单条数据删除。
sql
SET 'execution.runtime-mode' = 'batch';
DELETE FROM pk_table WHERE shop_id = 10000 AND user_id = 123456;
更新
Fluss 通过 UPDATE 语句以批处理模式为主键表启用数据更新。目前仅支持基于主键的单行更新。
sql
SET execution.runtime-mode = batch;
UPDATE pk_table SET total_amount = 2 WHERE shop_id = 10000 AND user_id = 123456;
Flink 读取
Fluss 支持使用 Apache Flink 的 SQL 和 Table API 进行流式和批量读取。执行以下 SQL 命令以在流式和批量模式之间切换执行模式:
sql
SET 'execution.runtime-mode' = 'streaming';
SET 'execution.runtime-mode' = 'batch';
流式读取
默认情况下,Streaming read 在首次启动时会在表上产生最新的快照,并继续读取最新的更改。Fluss 默认确保您的启动程序被正确处理,包含所有数据。流模式下的 Fluss Source 是无限的,就像一个永远不会结束的队列。
sql
SET 'execution.runtime-mode' = 'streaming';
SELECT * FROM my_table;
也可以在不读取快照数据的情况下进行流式读取,您可以使用 latest 扫描模式,该模式仅读取从最新偏移量开始的变化日志(或日志)。
sql
SELECT * FROM my_table /*+ OPTIONS('scan.startup.mode' = 'latest') */;
批量读取
Fluss 源支持限制读取
主键表和日志表都支持,便于预览表中的最新 N 条记录。
sql
-- 创建一张表并准备数据
CREATE TABLE log_table (
c_custkey INT NOT NULL,
c_name STRING NOT NULL,
c_address STRING NOT NULL,
c_nationkey INT NOT NULL,
c_phone STRING NOT NULL,
c_acctbal DECIMAL(15, 2) NOT NULL,
c_mktsegment STRING NOT NULL,
c_comment STRING NOT NULL
);
INSERT INTO log_table
VALUES (1, 'Customer1', 'IVhzIApeRb ot,c,E', 15, '25-989-741-2988', 711.56, 'BUILDING', 'comment1'),
(2, 'Customer2', 'XSTf4,NCwDVaWNe6tEgvwfmRchLXak', 13, '23-768-687-3665', 121.65, 'AUTOMOBILE', 'comment2'),
(3, 'Customer3', 'MG9kdTD2WBHm', 1, '11-719-748-3364', 7498.12, 'AUTOMOBILE', 'comment3');
-- 查询表
SET 'execution.runtime-mode' = 'batch';
SET 'sql-client.execution.result-mode' = 'tableau';
SELECT * FROM log_table LIMIT 10;
Fluss 源支持对主键表的点查询
允许您高效地检查特定记录。目前,此功能仅限于主键表。
sql
-- 创建一张主键表并准备数据
CREATE TABLE pk_table (
c_custkey INT NOT NULL,
c_name STRING NOT NULL,
c_address STRING NOT NULL,
c_nationkey INT NOT NULL,
c_phone STRING NOT NULL,
c_acctbal DECIMAL(15, 2) NOT NULL,
c_mktsegment STRING NOT NULL,
c_comment STRING NOT NULL,
PRIMARY KEY (c_custkey) NOT ENFORCED
);
INSERT INTO pk_table
VALUES (1, 'Customer1', 'IVhzIApeRb ot,c,E', 15, '25-989-741-2988', 711.56, 'BUILDING', 'comment1'),
(2, 'Customer2', 'XSTf4,NCwDVaWNe6tEgvwfmRchLXak', 13, '23-768-687-3665', 121.65, 'AUTOMOBILE', 'comment2'),
(3, 'Customer3', 'MG9kdTD2WBHm', 1, '11-719-748-3364', 7498.12, 'AUTOMOBILE', 'comment3');
-- 查询表
SET 'execution.runtime-mode' = 'batch';
SET 'sql-client.execution.result-mode' = 'tableau';
SELECT * FROM pk_table WHERE c_custkey = 1;
Fluss 源支持聚合
以批量模式对日志表进行下推计数聚合。这有助于预览日志表的总数。
sql
-- 在当前会话上下文中以批处理模式执行 Flink 作业
SET 'execution.runtime-mode' = 'batch';
SET 'sql-client.execution.result-mode' = 'tableau';
SELECT COUNT(*) FROM log_table;
读取选项
scan.startup.mode:扫描启动模式允许您指定数据消费的起始点。initial(默认):对于主键表,它首先消耗完整数据集,然后消耗增量数据。对于日志表,它从最早的偏移量开始消耗。earliest:对于主键表,它从最早的变更日志偏移量开始消费;对于日志表,它从最早的日志偏移量开始消费。latest:对于主键表,它从最新的变更日志偏移量开始消费;对于日志表,它从最新的日志偏移量开始消费。timestamp:对于主键表,它从指定的配置项scan.startup.timestamp定义的时间开始消费变更日志;对于日志表,它从对应指定时间的偏移量开始消费。
scan.partition.discovery.interval:Fluss 源在扫描分区表时发现新分区的时间间隔(毫秒)。默认值为 10 秒。非正值将禁用分区发现。
Flink Lookup Joins
Flink lookup joins 能够实现流数据与参考数据的实时高效丰富,这在许多实时分析和处理场景中是一个常见需求。
使用要求
- 使用主键表作为维度表,并且连接条件必须包含维度表的所有主键。
- Fluss 查找连接默认为异步模式以提高吞吐量。可以通过设置 SQL 提示
'lookup.async' = 'false'将查找连接的模式更改为同步模式。
使用案例
1: 创建两张表
sql
CREATE TABLE `fluss_catalog`.`my_db`.`orders` (
o_orderkey INT NOT NULL,
o_custkey INT NOT NULL,
o_orderstatus CHAR(1) NOT NULL,
o_totalprice DECIMAL(15, 2) NOT NULL,
o_orderdate DATE NOT NULL,
o_orderpriority CHAR(15) NOT NULL,
o_clerk CHAR(15) NOT NULL,
o_shippriority INT NOT NULL,
o_comment STRING NOT NULL,
PRIMARY KEY (o_orderkey) NOT ENFORCED
);
CREATE TABLE `fluss_catalog`.`my_db`.`customer` (
c_custkey INT NOT NULL,
c_name STRING NOT NULL,
c_address STRING NOT NULL,
c_nationkey INT NOT NULL,
c_phone CHAR(15) NOT NULL,
c_acctbal DECIMAL(15, 2) NOT NULL,
c_mktsegment CHAR(10) NOT NULL,
c_comment STRING NOT NULL,
PRIMARY KEY (c_custkey) NOT ENFORCED
);
2: 执行 Lookup Join
sql
USE CATALOG `fluss_catalog`;
USE my_db;
CREATE TEMPORARY TABLE lookup_join_sink (
order_key INT NOT NULL,
order_totalprice DECIMAL(15, 2) NOT NULL,
customer_name STRING NOT NULL,
customer_address STRING NOT NULL
) WITH ('connector' = 'blackhole');
-- 异步模式下的 Lookup Join
INSERT INTO lookup_join_sink
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM
(SELECT `orders`.*, proctime() AS ptime FROM `orders`) AS `o`
LEFT JOIN `customer`
FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
-- 同步模式下的 Lookup Join
INSERT INTO lookup_join_sink
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM
(SELECT `orders`.*, proctime() AS ptime FROM `orders`) AS `o`
LEFT JOIN `customer` /*+ OPTIONS('lookup.async' = 'false') */
FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
Lookup Options


大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...
共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术
- 面试指导、模拟面试
- 各类pdf文档下载、星球代码下载
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。
需要资料请加微信:D1435221412


end
