| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
| 1 年前 | ||
| 11 个月前 | ||
| 11 个月前 |
Apache Flink 是一个开源的分布式流处理框架,专为大规模数据处理设计,支持流处理和批处理一体化。它能够在无界和有界数据流上进行有状态计算,具有高吞吐量、低延迟的特性,广泛应用于实时数据处理场景。
本文档适用于Flink 2.0.x,Flink Connector JDBC 4.0.x版本。
相关依赖jar包下载地址:
使用说明
- gaussdb 数据库环境准备
gaussdb 数据库实例购买 以及库表创建。(以下为测试代码示例)
在gaussdb 数据库创建db test_db,创建schema player, 创建表 players
create database test_db;
use test_db;
create schema player;
CREATE TABLE player.players (
player_id INT NOT NULL,
team_id INT,
player_name VARCHAR(255),
height VARCHAR(255),
update_time timestamp,
PRIMARY KEY (player_id)
);
- flink sql 使用场景
- 做源表
场景介绍: Flink sql 读取gaussdb players数据后将数据写入printsink表
CREATE TABLE source (
player_id INT,
team_id INT,
player_name VARCHAR,
height VARCHAR,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:gaussdb://1.1.1.1:8000/test_db?currentSchema=player',
'username' = 'user',
'password' = '123456',
'table-name' = 'players');
CREATE TABLE printsink (
player_id INT,
team_id INT,
player_name VARCHAR,
height VARCHAR ,
update_time timestamp
) WITH (
'connector' = 'print'
);
insert into printsink select * from source;
执行成功:

- 做结果表
场景介绍: Flink sql 读取gaussdb players_source表的数据后将数据写入gaussdb players表
CREATE TABLE source (
player_id INT,
team_id INT,
player_name VARCHAR,
height VARCHAR,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:gaussdb://1.1.1.1:8000/test_db?currentSchema=player',
'username' = 'user',
'password' = '123456',
'table-name' = 'players_source');
CREATE TABLE sink (
player_id INT,
team_id INT,
player_name VARCHAR,
height VARCHAR,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:gaussdb://1.1.1.1:8000/test_db?currentSchema=player',
'username' = 'user',
'password' = '123456',
'table-name' = 'players');
insert into sink select * from source;
执行成功:

- 做维表
场景介绍: Flink sql 读取gaussdb players_info表数据并join gaussdb维表 players后 将数据写入players_detail 结果表
CREATE TABLE players (
player_id INT,
team_id INT,
player_name VARCHAR,
height VARCHAR,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:gaussdb://1.1.1.1:8000/test_db?currentSchema=player',
'username' = 'user',
'password' = '123456',
'table-name' = 'players');
CREATE TABLE players_info (
player_id INT,
player_name VARCHAR,
phone VARCHAR ,
address VARCHAR ,
email VARCHAR ,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:gaussdb://1.1.1.1:8000/test_db?currentSchema=player',
'username' = 'user',
'password' = '123456',
'table-name' = 'players_info');
CREATE TABLE players_detail (
player_id INT,
player_name VARCHAR,
phone VARCHAR ,
address VARCHAR ,
email VARCHAR ,
update_time timestamp,
PRIMARY KEY (player_id) NOT ENFORCED
) WITH (
'connector' = 'print'
);
insert into players_detail
select ps.player_id, ps.player_name, pi.phone ,pi.address ,pi.email ,pi.update_time from players_info as pi
inner join players as ps on ps.player_id = pi.player_id ;
执行成功:

注意事项:
- GaussDB驱动的选择: 依照Flink集群环境的JDK版本选择对应的驱动版本。
- 源码仓库地址: flink-connector-jdbc-gaussdb , flink-connector-jdbc-core
- 本文档的内容介绍是对 Flink Table API Connectors 方式的实现, Flink DataStream Connectors 方式请参考官网链接。