61a5db1b创建于 2025年7月29日历史提交
文件最后提交记录最后更新时间
1 年前
11 个月前
11 个月前
README

Apache Flink ‌ 是一个开源的分布式流处理框架,专为大规模数据处理设计,支持流处理和批处理一体化。它能够在无界和有界数据流上进行有状态计算,具有高吞吐量、低延迟的特性,广泛应用于实时数据处理场景。

本文档适用于Flink 2.0.x,Flink Connector JDBC 4.0.x版本。

相关依赖jar包下载地址:

使用说明

  1. 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)
);
  1. 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;

执行成功:
img.png

  • 做结果表
    场景介绍: 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;

执行成功:
img.png

  • 做维表
    场景介绍: 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 ;

执行成功:
img.png

注意事项:

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