flink-connector-gbase8a 使用说明
简介
Flink 应用程序可以通过连接器读取和写入各种外部系统。它支持多种格式,以便对数据进行编码和解码以匹配 Flink 的数据结构。用户可以通过flink-connector-gbase8a 完成和gbase8a的数据交换。该连接器共涉及三个jar包,其中flink-connector-gbase8a-V1.0.jar是一个精简jar,仅包含连接器的代码,不包含依赖包;其他两个包则为需要的依赖包。
flink-connector-gbase8a-V1.0.jar
connector-helper-V1.0.jar
gbase-connector-java-8.3.81.53-build55.5.7-bin.jar
用例说明
运行命令进入flink sql clients
./flink-1.16.0/bin/sql-client.sh embedded --jar flink-connector-jdbc-3.0.0.jar --jar connector-helper-
3.0.0-jar-with-dependencies.jar --jar gbase-connector-java-8.3.81.53-build55.5.7-bin.jar
用法-1
如何借助flink-connector-gbase8a在flink内建立8a对应的catalog?
在flink sql 客户端建立gbase8a catalog,可以用来显示8a的数据库表
CREATE CATALOG my_catalog WITH(
'type' = 'gbase8a',
'default-database' = 'test',
'username' = 'root',
'password' = 'admin',
'base-url' = 'jdbc:gbase://192.168.238.13:5258'
);
代表完成对8a catalog的注册,默认是default_catalog, 如果要切换catalog:
use catalog my_catalog;
此时完成catalog的切换,执行 show current catalog 则会显示当前的catalog为 my_catalog
通过show databases/show tables 可以完成8a侧数据库表的查看
用法-2
如何借助flink-connector-gbase8a完成一个简单的数据迁移作业?
--借助connector-datagen定义数据源
create table streamingSource(
id int,
name varchar(50),
age int
) WITH (
'connector' = 'datagen',
Option Required Default Type Description
url required N/A String
jdbc连接database用url,举例
jdbc:gbase://192.168.238.13:5258
vc-name optional 登录用户默认vc String 远程8a集群的虚拟集群的vcname
username required N/A String 8a用户名
password required N/A String 8a用户密码
driver optional com.gbase.jdbc.Driver String jdbc驱动命名
database-name required N/A String 8a数据库
table-name required N/A String 8a表名
connection.maxretry-timeout
optional Duration.ofSeconds(60) Duration jdbc连接最大的重试超时时长
Option Required Default Type Description
sink-mode optional NORMAL String
连接器有两种sink模式,一种
是NORMAL,一种是POC;两种
sink模式的使用场景不同,其
中POC模式底部走数据加载逻
辑,理论上经过调优之后有这
更高的TPS,仅支持insert
only;另外一种模式NORMAL
模式,底层走的jdbc批量数据
写入写出,可以在传输语义上
支持的更为精细,支持
insert,update,delete;
连接参数
sink参数
'fields.name.length'='10'
);
--定义一张动态表,用来映射8a里的某张实际的物理表
CREATE TABLE sink (
id int,
name varchar(50),
age int
) WITH (
'connector' = 'gbase8a',
'url' = 'jdbc:gbase://192.168.238.13:5258',
'vc-name'='vcname000005',
'database-name' = 'test',
'table-name' = 't2',
'sink-mode'='POC',
'username' = 'root',
'password' = 'admin',
'driver' = 'com.gbase.jdbc.Driver',
'sink.parallelism' = '1'
);
--通过运行下面sql语句拉起flink job
insert into sink select * from streamingSource;连接参数

sink参数


source参数

用法-3
如何借助flink-connector-gbase8a将远端mysql-cdc的变更数据同步到8a中?该场景要求connector具备
upsert语义,甚至要支持delete语义
关于如何配置mysql cdc,请参考
SET 'execution.checkpointing.interval' = '3s';
-- 在 Flink SQL中注册 MySQL 表 'orders'
CREATE TABLE orders_source (
order_id INT,
order_date TIMESTAMP(0),
customer_name STRING,
price DECIMAL(10, 5),
product_id INT,
order_status BOOLEAN,
PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'localhost',
'port' = '3306',
'username' = 'root',
'password' = '123456',
'database-name' = 'mydb',
'table-name' = 'orders');
-- 创建flink dynamic table,用来定义和8a中物理表的映射关系
CREATE TABLE orders_sink (
order_id INT,
order_date TIMESTAMP(0),
customer_name STRING,
price DECIMAL(10, 5),
product_id INT,
order_status BOOLEAN,
GBase8a type flink sql type note
BIGINT BIGINT
INT INT
SMALLINT SMALLINT
TINYINT TINYINT
FLOAT FLOAT
DOUBLE DOUBLE
DECIMAL DECIMAL
NUMERIC DECIMAL
CHAR CHAR
VARCHAR VARCHAR
TEXT VARCHAR
BLOB BYTES
LONGBLOB BYTES
DATE DATE
DATETIME TIMESTAMP(6) FLINK 无 datetime 类型
TIME TIME
TIMESTAMP TIMESTAMP
因为cdc需要下游的connector必须要支持upsert\delete的语义,因此在定义orders_sink的时候,此时
的sink-mode只能选择normal模式;
在建立flink动态表之前,需要在8a中预先建立实际的表,推荐8a中建立哈希分布表,并且是以order_id
字段进行hash分布的哈希分布表,此时的性能表现最优;
数据类型的映射
在flink中建表时,建议采用下面的映射关系进行建表,否则可能引起报错:
PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
'connector' = 'gbase8a',
'url' = 'jdbc:gbase://192.168.238.13:5258',
’sink-mode’='NORMAL',
'username' = 'root',
'password' = '123456',
'database-name' = 'mydb',
'table-name' = 'orders');
-- 将mysql change log 数据同步到gbase8a中
INSERT INTO orders_sink select * from orders_source;因为cdc需要下游的connector必须要支持upsert\delete的语义,因此在定义orders_sink的时候,此时
的sink-mode只能选择normal模式;
在建立flink动态表之前,需要在8a中预先建立实际的表,推荐8a中建立哈希分布表,并且是以order_id
字段进行hash分布的哈希分布表,此时的性能表现最优;
数据类型的映射
在flink中建表时,建议采用下面的映射关系进行建表,否则可能引起报错:

评论
热门帖子
- 12025-12-01浏览数:182759
- 22023-05-09浏览数:25052
- 42023-09-25浏览数:18521
- 52020-05-11浏览数:17526