数据异构 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

相关文章

windos server2022里的DFS配置的实现

《windosserver2022里的DFS配置的实现》DFS是WindowsServer操作系统提供的一种功能,用于在多台服务器上集中管理共享文件夹和文件的分布式存储解决方案,本文就来介绍一下wi... 目录什么是DFS?优势:应用场景:DFS配置步骤什么是DFS?DFS指的是分布式文件系统(Distr

NFS实现多服务器文件的共享的方法步骤

《NFS实现多服务器文件的共享的方法步骤》NFS允许网络中的计算机之间共享资源,客户端可以透明地读写远端NFS服务器上的文件,本文就来介绍一下NFS实现多服务器文件的共享的方法步骤,感兴趣的可以了解一... 目录一、简介二、部署1、准备1、服务端和客户端:安装nfs-utils2、服务端:创建共享目录3、服

SpringBoot使用Apache Tika检测敏感信息

《SpringBoot使用ApacheTika检测敏感信息》ApacheTika是一个功能强大的内容分析工具,它能够从多种文件格式中提取文本、元数据以及其他结构化信息,下面我们来看看如何使用Ap... 目录Tika 主要特性1. 多格式支持2. 自动文件类型检测3. 文本和元数据提取4. 支持 OCR(光学

Java内存泄漏问题的排查、优化与最佳实践

《Java内存泄漏问题的排查、优化与最佳实践》在Java开发中,内存泄漏是一个常见且令人头疼的问题,内存泄漏指的是程序在运行过程中,已经不再使用的对象没有被及时释放,从而导致内存占用不断增加,最终... 目录引言1. 什么是内存泄漏?常见的内存泄漏情况2. 如何排查 Java 中的内存泄漏?2.1 使用 J

JAVA系统中Spring Boot应用程序的配置文件application.yml使用详解

《JAVA系统中SpringBoot应用程序的配置文件application.yml使用详解》:本文主要介绍JAVA系统中SpringBoot应用程序的配置文件application.yml的... 目录文件路径文件内容解释1. Server 配置2. Spring 配置3. Logging 配置4. Ma

Python MySQL如何通过Binlog获取变更记录恢复数据

《PythonMySQL如何通过Binlog获取变更记录恢复数据》本文介绍了如何使用Python和pymysqlreplication库通过MySQL的二进制日志(Binlog)获取数据库的变更记录... 目录python mysql通过Binlog获取变更记录恢复数据1.安装pymysqlreplicat

Linux使用dd命令来复制和转换数据的操作方法

《Linux使用dd命令来复制和转换数据的操作方法》Linux中的dd命令是一个功能强大的数据复制和转换实用程序,它以较低级别运行,通常用于创建可启动的USB驱动器、克隆磁盘和生成随机数据等任务,本文... 目录简介功能和能力语法常用选项示例用法基础用法创建可启动www.chinasem.cn的 USB 驱动

Java 字符数组转字符串的常用方法

《Java字符数组转字符串的常用方法》文章总结了在Java中将字符数组转换为字符串的几种常用方法,包括使用String构造函数、String.valueOf()方法、StringBuilder以及A... 目录1. 使用String构造函数1.1 基本转换方法1.2 注意事项2. 使用String.valu

C#使用yield关键字实现提升迭代性能与效率

《C#使用yield关键字实现提升迭代性能与效率》yield关键字在C#中简化了数据迭代的方式,实现了按需生成数据,自动维护迭代状态,本文主要来聊聊如何使用yield关键字实现提升迭代性能与效率,感兴... 目录前言传统迭代和yield迭代方式对比yield延迟加载按需获取数据yield break显式示迭

Python实现高效地读写大型文件

《Python实现高效地读写大型文件》Python如何读写的是大型文件,有没有什么方法来提高效率呢,这篇文章就来和大家聊聊如何在Python中高效地读写大型文件,需要的可以了解下... 目录一、逐行读取大型文件二、分块读取大型文件三、使用 mmap 模块进行内存映射文件操作(适用于大文件)四、使用 pand