GBase 8a
其他
文章

GBase 8a与flinkSQL联调

发表于2024-05-30 17:34:57167次浏览0个评论

目的:实现8a中的t1表,关联kafka数据,插入到8a中的t2表。当kafka数据变化时,t2表发生变化。


一、部署flink 、kafka 
二、创建表
1、在8a中创建
t1和t2表。
create table test.t1(id int,name varchar(10) ,primary key(id));
create table test .t2 like t1;
insert into test.t1 values(1,'a'),(2,'b');

2、在 Flink SQL 中注册对应表。 ./sql-client.sh
CREATE TABLE t1 (
 id int,
 name STRING
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://172.16.6.98:5258/test?useLocalSessionState=true',
  'driver' = 'com.mysql.jdbc.Driver',
  'table-name' = 't1',
  'username' = 'gbase',
  'password' = 'gbase20110531',
  'sink.buffer-flush.max-rows' = '200',  -- 批量输出的条数
'sink.buffer-flush.interval' = '1s'    -- 批量输出的间隔
);


CREATE TABLE t2 (
 id int,
 name STRING
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://172.16.6.98:5258/test?useLocalSessionState=true',
  'driver' = 'com.mysql.jdbc.Driver',
  'table-name' = 't2',
  'username' = 'gbase',
  'password' = 'gbase20110531'
);

CREATE TABLE kt2 (
id int
) WITH (
 'connector' = 'kafka',
 'topic' = 't1',
 'properties.bootstrap.servers' = '172.16.6.98:9092',
 'properties.group.id' = 'testGroup',
 'scan.startup.mode' = 'earliest-offset',
 'format' = 'csv'
);


insert into t2
select t1.id,t1.name 
from t1 inner join kt1 
on t1.id=kt1.id;

3、向kafka 生产者中输入1;
4、t2表数据会发生变化
 

 

评论

登录后才可以发表评论