改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Fluss初窥四与Flinksql集成使用

ckckck2025年1月11日22 浏览

Flink DDL

创建 Fluss Catalog

sql 复制代码
CREATE CATALOG fluss_catalog WITH (
  'type' = 'fluss',
  'bootstrap.servers' = 'fluss-server-1:9123'
);

USE CATALOG fluss_catalog;

以下属性在使用 Fluss 目录时可以设置:

Fluss Catalog Properties

创建数据库

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”参数支持的选项如下:

With Parameters

表配置参数:

Table Configuration Parameters

客户端配置参数:

Client Configuration Parameters

使用 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 主键表可以接受所有类型的消息(INSERTUPDATE_BEFOREUPDATE_AFTERDELETE),而 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

Lookup Options
Lookup Options

大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...

共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:

  1. 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
  2. 免费答疑解惑、交流技术
  3. 面试指导、模拟面试
  4. 各类pdf文档下载、星球代码下载
  5. 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。

另外说明加入星球后支持三天无理由退款,不满意无条件随时退。

需要资料请加微信:D1435221412

WeChat QR Code
WeChat QR Code

end