数据异构 Canal-Spring-Boot-Starter的技术实现

2024-05-04 01:08

本文主要是介绍数据异构 Canal-Spring-Boot-Starter的技术实现,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

Canal-Spring-Boot-Starter 使用

1、在spring boot 项目配置文件 application.yml内增加以下内容


spring:canal:instances:example:                  # 拉取 example 目标的数据host: 192.168.10.179    # canal 所在机器的ipport: 11111             # canal 默认暴露端口user-name: canal        # canal 用户名password: canal         # canal 密码batch-size: 600         # canal 每次拉取的数据条数retry-count: 5          # 重试次数,如果重试5次后,仍无法连接,则断开cluster-enabled: false  # 是否开启集群zookeeper-address:      # zookeeper 地址(开启集群的情况下生效), 例: 192.168.0.1:2181,192.168.0.2:2181,192.168.0.3:2181acquire-interval: 1000  # 未拉取到消息情况下,获取消息的时间间隔毫秒值subscribe: .*\\..*      # 默认情况下拉取所有库、所有表
prod:example: exampledatabase: books

2、在spring boot 项目中的代码使用实例

import com.alibaba.otter.canal.protocol.CanalEntry;
import com.duxinglangzi.canal.starter.annotation.CanalInsertListener;
import com.duxinglangzi.canal.starter.annotation.CanalListener;
import com.duxinglangzi.canal.starter.annotation.CanalUpdateListener;
import com.duxinglangzi.canal.starter.annotation.EnableCanalListener;
import com.duxinglangzi.canal.starter.mode.CanalMessage;
import org.springframework.stereotype.Service;import java.util.stream.Collectors;/*** @author wuqiong 2022/4/12* @description*/
@EnableCanalListener
@Service
public class CanalListenerTest {/*** 必须在类上 使用 EnableCanalListener 注解才能开启 canal listener** 目前 Listener 方法的参数必须为 com.duxinglangzi.canal.starter.mode.CanalMessage* 程序在启动过程中会做检查*//*** 监控更新操作* 支持动态参数配置,配置项需在 yml 或 properties 进行配置* 目标是 ${prod.example} 的  ${prod.database} 库  users表*/@CanalUpdateListener(destination = "${prod.example}", database = "${prod.database}", table = {"users"})public void listenerExampleBooksUsers(CanalMessage message) {printChange("listenerExampleBooksUsers", message);}/*** 监控更新操作 ,目标是 example的  books库  users表*/@CanalInsertListener(destination = "example", database = "books", table = {"users"})public void listenerExampleBooksUser(CanalMessage message) {printChange("listenerExampleBooksUsers", message);}/*** 监控更新操作 ,目标是 example的  books库  books表*/@CanalUpdateListener(destination = "example", database = "books", table = {"books"})public void listenerExampleBooksBooks(CanalMessage message) {printChange("listenerExampleBooksBooks", message);}/*** 监控更新操作 ,目标是 example的  books库的所有表*/@CanalListener(destination = "example", database = "books", eventType = CanalEntry.EventType.UPDATE)public void listenerExampleBooksAll(CanalMessage message) {printChange("listenerExampleBooksAll", message);}/*** 监控更新操作 ,目标是 example的  所有库的所有表*/@CanalListener(destination = "example", eventType = CanalEntry.EventType.UPDATE)public void listenerExampleAll(CanalMessage message) {printChange("listenerExampleAll", message);}/*** 监控更新、删除、新增操作 ,所有配置的目标下的所有库的所有表*/@CanalListener(eventType = {CanalEntry.EventType.UPDATE, CanalEntry.EventType.INSERT, CanalEntry.EventType.DELETE})public void listenerAllDml(CanalMessage message) {printChange("listenerAllDml", message);}public void printChange(String method, CanalMessage message) {CanalEntry.EventType eventType = message.getEventType();CanalEntry.RowData rowData = message.getRowData();System.out.println(" >>>>>>>>>>>>>[当前数据库: "+message.getDataBaseName()+" ," +"数据库表名: " + message.getTableName() + " , " +"方法: " + method );if (eventType == CanalEntry.EventType.DELETE) {rowData.getBeforeColumnsList().stream().collect(Collectors.toList()).forEach(ele -> {System.out.println("[方法: " + method + " ,  delete 语句 ] --->> 字段名: " + ele.getName() + ", 删除的值为: " + ele.getValue());});}if (eventType == CanalEntry.EventType.INSERT) {rowData.getAfterColumnsList().stream().collect(Collectors.toList()).forEach(ele -> {System.out.println("[方法: " + method + " ,insert 语句 ] --->> 字段名: " + ele.getName() + ", 新增的值为: " + ele.getValue());});}if (eventType == CanalEntry.EventType.UPDATE) {for (int i = 0; i < rowData.getAfterColumnsList().size(); i++) {CanalEntry.Column afterColumn = rowData.getAfterColumnsList().get(i);CanalEntry.Column beforeColumn = rowData.getBeforeColumnsList().get(i);System.out.println("[方法: " + method + " , update 语句 ] -->> 字段名," + afterColumn.getName() +" , 是否修改: " + afterColumn.getUpdated() +" , 修改前的值: " + beforeColumn.getValue() +" , 修改后的值: " + afterColumn.getValue());}}}}

以上展示了在 spring boot 项目中 canal starter 的基本使用

3、源码地址

对于 canal-spring-boot-starter 源代码为楼主自己封装, github地址: https://github.com/duxinglangzi/canal-spring-boot-starter
另附国内 gitee 地址: https://gitee.com/duxinglangzi/canal-spring-boot-starter

如果有需要的同学,可以自行下载源码进行修改和自定义封装

这篇关于数据异构 Canal-Spring-Boot-Starter的技术实现的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/958042

相关文章

springboot集成easypoi导出word换行处理过程

《springboot集成easypoi导出word换行处理过程》SpringBoot集成Easypoi导出Word时,换行符n失效显示为空格,解决方法包括生成段落或替换模板中n为回车,同时需确... 目录项目场景问题描述解决方案第一种:生成段落的方式第二种:替换模板的情况,换行符替换成回车总结项目场景s

SpringBoot集成redisson实现延时队列教程

《SpringBoot集成redisson实现延时队列教程》文章介绍了使用Redisson实现延迟队列的完整步骤,包括依赖导入、Redis配置、工具类封装、业务枚举定义、执行器实现、Bean创建、消费... 目录1、先给项目导入Redisson依赖2、配置redis3、创建 RedissonConfig 配

SpringBoot中@Value注入静态变量方式

《SpringBoot中@Value注入静态变量方式》SpringBoot中静态变量无法直接用@Value注入,需通过setter方法,@Value(${})从属性文件获取值,@Value(#{})用... 目录项目场景解决方案注解说明1、@Value("${}")使用示例2、@Value("#{}"php

SpringBoot分段处理List集合多线程批量插入数据方式

《SpringBoot分段处理List集合多线程批量插入数据方式》文章介绍如何处理大数据量List批量插入数据库的优化方案:通过拆分List并分配独立线程处理,结合Spring线程池与异步方法提升效率... 目录项目场景解决方案1.实体类2.Mapper3.spring容器注入线程池bejsan对象4.创建

线上Java OOM问题定位与解决方案超详细解析

《线上JavaOOM问题定位与解决方案超详细解析》OOM是JVM抛出的错误,表示内存分配失败,:本文主要介绍线上JavaOOM问题定位与解决方案的相关资料,文中通过代码介绍的非常详细,需要的朋... 目录一、OOM问题核心认知1.1 OOM定义与技术定位1.2 OOM常见类型及技术特征二、OOM问题定位工具

PHP轻松处理千万行数据的方法详解

《PHP轻松处理千万行数据的方法详解》说到处理大数据集,PHP通常不是第一个想到的语言,但如果你曾经需要处理数百万行数据而不让服务器崩溃或内存耗尽,你就会知道PHP用对了工具有多强大,下面小编就... 目录问题的本质php 中的数据流处理:为什么必不可少生成器:内存高效的迭代方式流量控制:避免系统过载一次性

Python的Darts库实现时间序列预测

《Python的Darts库实现时间序列预测》Darts一个集统计、机器学习与深度学习模型于一体的Python时间序列预测库,本文主要介绍了Python的Darts库实现时间序列预测,感兴趣的可以了解... 目录目录一、什么是 Darts?二、安装与基本配置安装 Darts导入基础模块三、时间序列数据结构与

基于 Cursor 开发 Spring Boot 项目详细攻略

《基于Cursor开发SpringBoot项目详细攻略》Cursor是集成GPT4、Claude3.5等LLM的VSCode类AI编程工具,支持SpringBoot项目开发全流程,涵盖环境配... 目录cursor是什么?基于 Cursor 开发 Spring Boot 项目完整指南1. 环境准备2. 创建

Python使用FastAPI实现大文件分片上传与断点续传功能

《Python使用FastAPI实现大文件分片上传与断点续传功能》大文件直传常遇到超时、网络抖动失败、失败后只能重传的问题,分片上传+断点续传可以把大文件拆成若干小块逐个上传,并在中断后从已完成分片继... 目录一、接口设计二、服务端实现(FastAPI)2.1 运行环境2.2 目录结构建议2.3 serv

C#实现千万数据秒级导入的代码

《C#实现千万数据秒级导入的代码》在实际开发中excel导入很常见,现代社会中很容易遇到大数据处理业务,所以本文我就给大家分享一下千万数据秒级导入怎么实现,文中有详细的代码示例供大家参考,需要的朋友可... 目录前言一、数据存储二、处理逻辑优化前代码处理逻辑优化后的代码总结前言在实际开发中excel导入很