聚合搜索平台-SpringBoot整合Canal

官方文档:https://github.com/alibaba/canal/wiki/QuickStart

1.准备工作

开启MySQL的Binglog功能配置 binlog-format 为 ROW 模式,my.cnf 中配置如下

nginx
复制代码
[mysqld] log-bin=mysql-bin # 开启 binlog binlog-format=ROW # 选择 ROW 模式 server_id=1 # 配置 MySQL replaction 需要定义,不要和 canal 的 slaveId 重复

创建用户

sql
复制代码
CREATE USER canal IDENTIFIED BY 'canal'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%'; -- GRANT ALL PRIVILEGES ON *.* TO 'canal'@'%' ; FLUSH PRIVILEGES;

如果是mysql8要运行以下命令,

这是由于mysql8与以前版本的mysql密码算法不同的原因

sql
复制代码
ALTER USER 'canal'@'%' IDENTIFIED WITH mysql_native_password by 'canal'; ALTER USER 'canal'@'%' IDENTIFIED BY 'canal' PASSWORD EXPIRE NEVER; FLUSH PRIVILEGES;

根据快速开始,选择下载版本,1.1.8

修改配置

conf/example/instance.properties

properties
复制代码
## mysql serverId canal.instance.mysql.slaveId = 1234 #position info,需要改成自己的数据库信息 canal.instance.master.address = 127.0.0.1:3306 canal.instance.master.journal.name = canal.instance.master.position = canal.instance.master.timestamp = #canal.instance.standby.address = #canal.instance.standby.journal.name = #canal.instance.standby.position = #canal.instance.standby.timestamp = #username/password,需要改成自己的数据库信息 canal.instance.dbUsername = canal canal.instance.dbPassword = canal canal.instance.defaultDatabaseName = canal.instance.connectionCharset = UTF-8 #table regex canal.instance.filter.regex = .\*\\\\..\*

注意:只是修改配置文件中的部分配置,不需要全部覆盖,否则会导致无法连接

依赖

选择对应版本的依赖

xml
复制代码
<dependency> <groupId>com.alibaba.otter</groupId> <artifactId>canal.client</artifactId> <version>1.1.8</version> </dependency> <dependency> <groupId>com.alibaba.otter</groupId> <artifactId>canal.protocol</artifactId> <version>1.1.8</version> </dependency>

1.1 简单示例

https://github.com/alibaba/canal/wiki/ClientExample

java
复制代码
import com.alibaba.otter.canal.client.CanalConnector; import com.alibaba.otter.canal.client.CanalConnectors; import com.alibaba.otter.canal.common.utils.AddressUtils; import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.Message; import java.net.InetSocketAddress; import java.util.List; public class SimpleCanalClientExample { public static void main(String args[]) { // 创建链接 CanalConnector connector = CanalConnectors.newSingleConnector(new InetSocketAddress(AddressUtils.getHostIp(), 11111), "example", "", ""); int batchSize = 1000; int emptyCount = 0; try { connector.connect(); connector.subscribe(".*\\..*"); connector.rollback(); int totalEmptyCount = 120; while (emptyCount < totalEmptyCount) { Message message = connector.getWithoutAck(batchSize);// 获取指定数量的数据 long batchId = message.getId(); int size = message.getEntries().size(); if (batchId == -1 || size == 0) { emptyCount++; System.out.println("empty count : " + emptyCount); try { Thread.sleep(1000); } catch (InterruptedException e) { } } else { emptyCount = 0; // System.out.printf("message[batchId=%s,size=%s] \n", batchId, size); List<CanalEntry.Entry> entries = message.getEntries(); printEntry(entries); } connector.ack(batchId); // 提交确认 // connector.rollback(batchId); // 处理失败, 回滚数据 } System.out.println("empty too many times, exit"); } finally { connector.disconnect(); } } private static void printEntry(List<CanalEntry.Entry> entrys) { for (CanalEntry.Entry entry : entrys) { if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN || entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) { continue; } CanalEntry.RowChange rowChage = null; try { rowChage = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e); } CanalEntry.EventType eventType = rowChage.getEventType(); System.out.println(String.format("================&gt; binlog[%s:%s] , name[%s,%s] , eventType : %s", entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType)); for (CanalEntry.RowData rowData : rowChage.getRowDatasList()) { if (eventType == CanalEntry.EventType.DELETE) { printColumn(rowData.getBeforeColumnsList()); } else if (eventType == CanalEntry.EventType.INSERT) { printColumn(rowData.getAfterColumnsList()); } else { System.out.println("-------&gt; before"); List<CanalEntry.Column> beforeColumnsList = rowData.getBeforeColumnsList(); printColumn(beforeColumnsList); System.out.println("-------&gt; after"); printColumn(rowData.getAfterColumnsList()); } } } } private static void printColumn(List<CanalEntry.Column> columns) { for (CanalEntry.Column column : columns) { System.out.println(column.getName() + " : " + column.getValue() + " update=" + column.getUpdated()); } } }

2.与springboot集成

这里只实现了一个数据库的监控,监控所有数据库应该也时差不多,但是要注意数据库之间的实体类和Es文档映射的问题

可能只试用于 MySQL ---》Es

2.1. Es文档类

java
复制代码
package com.sakura.demo.datasync.modal.document; import lombok.Data; import org.springframework.data.annotation.Id; import org.springframework.data.elasticsearch.annotations.*; import java.math.BigDecimal; import java.util.Date; @Document(indexName = "product", createIndex = true) @Data public class ProductDocument { @Id private Long id; @MultiField( mainField = @Field( type = FieldType.Text, analyzer = "ik_max_word", // 索引时使用细粒度分词 searchAnalyzer = "ik_smart" // 搜索时使用智能分词 ), otherFields = { @InnerField( suffix = "keyword", type = FieldType.Keyword, ignoreAbove = 256 // 超过256字符的keyword会被忽略 ) } ) private String name; @MultiField( mainField = @Field( type = FieldType.Text, analyzer = "ik_max_word", searchAnalyzer = "ik_smart" ), otherFields = { @InnerField(suffix = "keyword", type = FieldType.Keyword, ignoreAbove = 256) } ) private String description; @Field(type = FieldType.Double) private BigDecimal price; // 忽略映射 @Field(index = false, store = true, type = FieldType.Date, format = {}, pattern = DATE_TIME_PATTERN) private Date updateTime; private static final String DATE_TIME_PATTERN = "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'"; }

2.2. 编写配置类

秉承着 约定>配置>编码,因此我们要实现一个配置类,来配置canalClient

java
复制代码
/** * canal配置 */ @Component @Data @ConfigurationProperties("canal") public class CanalProperties { /** * 端口 */ private Integer port; /** * 描述 */ private String destination; /** * 用户名 */ private String username; /** * 密码 */ private String password; /** * 批量大小 */ private Integer batchSize; /** * 过滤 */ private String filter; /** * 实体类包名前缀 */ private String basePackage; /** * 实体类后缀 */ private String clazzSuffix; }

Spring配置文件

yaml
复制代码
# canal 配置 canal: port: 11111 batch-size: 1000 destination: example # https://github.com/alibaba/canal/wiki/AdminGuide canal.instance.filter.regex配置 filter: "my_db\\..*" username: password: # base-package 和 clazz-suffix 要与你文档标注的类一致 # 例如你有一个 com.sakura.demo.datasync.modal.document.ProductDocument 类 # 那么 base-package:com.sakura.demo.datasync.modal.document # clazz-suffix: Document base-package: com.sakura.demo.datasync.modal.document clazz-suffix: Document

config类

java
复制代码
/** * Canal配置 */ @Slf4j @Configuration public class CanalConfig { @Bean(destroyMethod = "disconnect") public CanalConnector canalConnector(CanalProperties properties) { log.info("Canal配置初始化"); return CanalConnectors.newSingleConnector(new InetSocketAddress(AddressUtils.getHostIp(), properties.getPort()), properties.getDestination(), properties.getUsername(), properties.getPassword()); } }

特别说明:

  1. base-package 和 clazz-suffix 要与你文档标注的类一致。

  2. 例如你有一个 com.sakura.demo.datasync.modal.document.ProductDocument 类

  3. 那么 base-package:com.sakura.demo.datasync.modal.document

  4. clazz-suffix: Document

2.3. Canal监听数据服务

为了方便查看,差分成下面几个代码模块

java
复制代码
import com.alibaba.otter.canal.client.CanalConnector; import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.Message; import com.sakura.demo.datasync.modal.properties.CanalProperties; import com.sakura.demo.datasync.utils.StringUtils; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.DisposableBean; import org.springframework.data.elasticsearch.core.ElasticsearchRestTemplate; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import javax.annotation.Resource; import java.lang.reflect.Field; import java.math.BigDecimal; import java.text.ParseException; import java.text.SimpleDateFormat; import java.util.Date; import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @Slf4j @Component public class CanalDataSyncService implements DisposableBean { @Resource private CanalConnector canalConnector; @Resource private CanalProperties canalProperties; @Resource private ElasticsearchRestTemplate restTemplate; private ExecutorService executorService; @PostConstruct public void init() { // 初始化线程池 executorService = Executors.newSingleThreadExecutor(); // 提交任务到线程池 executorService.submit(this::process); log.info("Canal 开始监控数据改变"); } private void process() { canalConnector.connect(); canalConnector.subscribe(canalProperties.getFilter()); while (true) { Message message = canalConnector.get(100); List<CanalEntry.Entry> entries = message.getEntries(); // 处理数据 handlerEntries(entries); } } @Override public void destroy() throws Exception { // 关闭CanalConnector连接 canalConnector.disconnect(); log.info("CanalConnector已断开连接"); // 关闭线程池 if (executorService != null) { executorService.shutdown(); try { if (!executorService.awaitTermination(60, TimeUnit.SECONDS)) { executorService.shutdownNow(); } } catch (InterruptedException e) { executorService.shutdownNow(); } log.info("线程池已关闭"); } } }

handlerEntries方法

java
复制代码
private void handlerEntries(List<CanalEntry.Entry> entrys) { for (CanalEntry.Entry entry : entrys) { if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN || entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) { continue; } CanalEntry.RowChange rowChage = null; try { rowChage = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e); } CanalEntry.EventType eventType = rowChage.getEventType(); log.info("================> binlog[{}:{}] , name[{},{}] , eventType : {}", entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType); String tableName = entry.getHeader().getTableName(); // 表面对应实体类名称 // 这里利用反射拿到对应的实体类 Class<?> clazz = null; Object syncObj = null; try { String clazzName = StringUtils.toPascalCase(tableName); clazz = Class.forName(canalProperties.getBasePackage() + "." + clazzName + canalProperties.getClazzSuffix()); log.info("获取到的实体类:{}", clazz); syncObj = clazz.getDeclaredConstructor().newInstance(); } catch (Exception e) { log.error("获取实体类失败,表名:{}", tableName, e); continue; // 跳过当前条目,继续处理下一个 } for (CanalEntry.RowData rowData : rowChage.getRowDatasList()) { handleRowData(rowData, eventType, clazz, syncObj); } } }

handleRowData方法

java
复制代码
/** * 处理行数据 * @param rowData 行数据 * @param eventType 事件类型 * @param clazz 反射类 * @param syncObj 同步实体对象 */ private void handleRowData(CanalEntry.RowData rowData, CanalEntry.EventType eventType, Class<?> clazz, Object syncObj) { if (eventType == CanalEntry.EventType.DELETE) { // 处理删除操作 List<CanalEntry.Column> beforeColumnsList = rowData.getBeforeColumnsList(); // 只需要拿到对应的id就好 handlerColumn(beforeColumnsList, clazz, syncObj); restTemplate.delete(syncObj); } else { // 处理插入/更新操作 List<CanalEntry.Column> afterColumnsList = rowData.getAfterColumnsList(); handlerColumn(afterColumnsList, clazz, syncObj); restTemplate.save(syncObj); } }

handlerColumn方法

java
复制代码
/** * 处理列数据 * @param columns 列数据 * @param clazz 反射类 * @param syncObj 同步实体对象 */ private void handlerColumn(List<CanalEntry.Column> columns, Class<?> clazz, Object syncObj) { for (CanalEntry.Column column : columns) { log.info("{} : {} update={}", column.getName(), column.getValue(), column.getUpdated()); try { Field field = clazz.getDeclaredField(StringUtils.toCamelCase(column.getName())); field.setAccessible(true); // 根据字段类型设置值 setFieldValue(field, syncObj, column.getValue()); } catch (Exception e) { log.error("设置实体类属性失败,字段名:{}", column.getName(), e); } } } /** * 处理字段映射 * @param field 字段 * @param obj 实体对象 * @param value 字段值 * @throws IllegalAccessException */ private static void setFieldValue(Field field, Object obj, String value) throws IllegalAccessException { Class<?> fieldType = field.getType(); if (fieldType == String.class) { field.set(obj, value); } else if (fieldType == int.class || fieldType == Integer.class) { field.set(obj, Integer.parseInt(value)); } else if (fieldType == long.class || fieldType == Long.class) { field.set(obj, Long.parseLong(value)); } else if (fieldType == double.class || fieldType == Double.class) { field.set(obj, Double.parseDouble(value)); } else if (fieldType == boolean.class || fieldType == Boolean.class) { field.set(obj, Boolean.parseBoolean(value)); } else if (fieldType == Date.class) { SimpleDateFormat dateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); try { field.set(obj, dateFormat.parse(value)); } catch (ParseException e) { log.error("无法解析日期值: {}", value, e); throw new IllegalArgumentException("无法解析日期值: " + value, e); } } else if (fieldType == BigDecimal.class) { field.set(obj, new BigDecimal(value)); } else { throw new IllegalArgumentException("Unsupported field type: " + fieldType); } }

Gitee地址:https://gitee.com/kk-2049/learn-dem Github地址:https://github.com/modakai/learn-demo

有些地方可能写的不太好,望见谅和指正

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP