使用flink 消费kafka中的canal-json数据到GBase8a 范例
一、逻辑说明
由于GBase 支持直接对接kafka ,该项目方案在前期讨论 共2讨论了2种,综合考虑,使
用方案2 实现难度更低,对产品稳定性影响较小。
1、GBase 直接对接kafka 进行消费。
—— 产品已经支持 人保、OGG 等3类 json 数据消费,目前不支持天津公积金的
json 格式,需要二次开发。开发基于产品内置功能修改,开发、测试的周期整体较长,属
于定制型开发,暂时不使用该方案。
2、使用flink 进行处理
—— 实现 flink connector ,通过flink 对接kafka与GBase,将kafka 的实时数据消
费到GBase中。
二、注意事项
1、由于 验证环境是JDK1.8版本 ,部分新版本kafka 基于JDK1.7编译的,不试用。
2、flink 会发送参数 allowAutoTopicCreation,在低版本的kafka上 并不支持 ,
MetadataRequest 在kafka的 v4 上才支持。。
3、本次测试的flink 版本是1.8 ,kafka版本是 2.13-2.8.2 。 现场如果已经使用flink 对
接kafka 的实时数据消费,理论上版本号是满足的。
三、 导入flink connector 包
1、 将包解压后,放入 flink的 lib 目录下。重启flink 服务即可。
或
2、启动 flinkSQL 客户端时 指定这3个jar包,参考:
# 启动示例
路径/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、数据样例
新增数据
{"data":
[{"id":"2","device_id":"2","metric_name":"22","metric_value":"2.0
","event_time":"2026-07-13
00:00:00","note":"33","GTID":"13"}],"database":"test","es":178392
3670000,"gtid":"","id":187,"isDdl":false,"mysqlType":
{"id":"decimal(10,0)","device_id":"varchar(32)","metric_name":"va
rchar(32)","metric_value":"decimal(10,2)","event_time":"datetime"
,"note":"varchar(128)","GTID":"bigint
unsigned"},"old":null,"pkNames":["id"],"sql":"","sqlType":
{"id":3,"device_id":12,"metric_name":12,"metric_value":3,"event_t
ime":93,"note":12,"GTID":-5},"table":"oracle_dup_canal_probe","ts
":1783923830267,"type":"INSERT"}
修改数据
{"data":
[{"id":"2","device_id":"2","metric_name":"22","metric_value":"2.0
","event_time":"2026-07-13
00:00:00","note":"33","GTID":"14"}],"database":"test","es":178392
3685000,"gtid":"","id":189,"isDdl":false,"mysqlType":
{"id":"decimal(10,0)","device_id":"varchar(32)","metric_name":"va
rchar(32)","metric_value":"decimal(10,2)","event_time":"datetime"
,"note":"varchar(128)","GTID":"bigint unsigned"},"old":
[{"GTID":"13"}],"pkNames":["id"],"sql":"","sqlType":
{"id":3,"device_id":12,"metric_name":12,"metric_value":3,"event_t
ime":93,"note":12,"GTID":-5},"table":"oracle_dup_canal_probe","ts
":1783923844881,"type":"UPDATE"}
删除数据
{"data":
[{"id":"2","device_id":"2","metric_name":"22","metric_value":"2.0
","event_time":"2026-07-13
00:00:00","note":"33","GTID":"14"}],"database":"test","es":178392
3685000,"gtid":"","id":189,"isDdl":false,"mysqlType":
{"id":"decimal(10,0)","device_id":"varchar(32)","metric_name":"va
rchar(32)","metric_value":"decimal(10,2)","event_time":"datetime"
,"note":"varchar(128)","GTID":"bigint
unsigned"},"old":null,"pkNames":["id"],"sql":"","sqlType":
{"id":3,"device_id":12,"metric_name":12,"metric_value":3,"event_t
ime":93,"note":12,"GTID":-5},"table":"oracle_dup_canal_probe","ts
":1783923844881,"type":"DELETE"}
2、 搭建 flink 、kafka 环境 。
3、 在 GBase 8a 中创建目标表,用于验证数据是否能正常消费
4、登录flink sql 客户端,创建 映射kafka的source 映射表,创建 GBase 的sink 映射表,创
建数据同步任务。
-- 创建kafka source 映射表
CREATE TABLE kafka_canal_source (
id DECIMAL(10,0),
device_id VARCHAR(32),
metric_name VARCHAR(32),
metric_value DECIMAL(10,2),
event_time TIMESTAMP(0),
note VARCHAR(128),
GTID BIGINT,
kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp',
-- 业务主键id,Flink Kafka源仅支持逻辑主键 NOT ENFORCED
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'kafka',
'topic' = 'testload',
'properties.bootstrap.servers' =
'172.16.6.103:9092,172.16.6.104:9092,172.16.6.105:9092',
'properties.group.id' = 'flink_canal_group',
'scan.startup.mode' = 'latest-offset',
'format' = 'canal-json',
19'canal-json.ignore-parse-errors' = 'true',
-- 可选:开启批量读取优化
'properties.fetch.min.bytes' = '1024',
'properties.fetch.max.wait.ms' = '500'
);
-- 创建 GBase 的sink 表
CREATE TABLE gbase_canal_sink
(
id DECIMAL(10,0) ,
device_id VARCHAR(32) ,
metric_name VARCHAR(32),
metric_value DECIMAL(10,2),
event_time TIMESTAMP,
note VARCHAR(128),
GTID BIGINT,
PRIMARY KEY (id) NOT ENFORCED
)
WITH (
'connector' = 'gbase8a',
'url' = 'jdbc:gbase://172.16.6.100:5258',
'host.list' = '172.16.6.101,172.16.6.102',
'database-name' = 'test',
'sink.mode' = 'NORMAL',
'table-name' = 'oracle_dup_canal_probe',
'vc-name'='vcname000001',
'driver' = 'com.gbase.jdbc.Driver',
'password' = 'gbase20110531',
'channel.buffer.batch.size'='20000',
'username' = 'gbase',
'sink.parallelism' = '1',
'sink.buffer-flush.max-rows'='500',
'sink.buffer-flush.interval'='1s',
'load.suffix' = 'DATETIME FORMAT [%Y-%m-%d %H:%i:%s.%f]'
) ;
-- 提交流式同步SQL
INSERT INTO gbase_canal_sink
SELECT id, device_id, metric_name, metric_value, event_time,
note, cast(GTID as bigint)
FROM kafka_canal_source;
5、验证方式
确保上面正常运行,没有报错。
A 、链接kafka 对应的topic ,使用生产者。
B、 链接kafka 对应的topic, 使用消费者。
C、生产者: 传入 第一个 insert 对应数据, 在消费者这能正常显示,说明kafka正常。
此时链接gbase ,查询该表(test.oracle_dup_canal_probe)数据是否正常。
D、生产者: 传入 第二个 update 对应数据, 在消费者这能正常显示,说明kafka正
常。此时链接gbase ,查询该表(test.oracle_dup_canal_probe)数据是否被更新。
E、生产者: 传入 第一个 insert 对应数据, 在消费者这能正常显示,说明kafka正常。
此时链接gbase ,查询该表(test.oracle_dup_canal_probe)数据是否消失。
评论
热门帖子
- 12025-12-01浏览数:182759
- 22023-05-09浏览数:25044
- 42023-09-25浏览数:18519
- 52020-05-11浏览数:17526