SaaS本地部署canal监听binlog迁移信创库技术方案
一、背景saas业务系统经常通过mysql binlog读取某些表的业务变更然后异步发送到MQ中间件为提升性能传输协议是protobuf而不是json再被各个业务模块进行聚合处理业务架构如下图所示如果把saas软件做本地化私有部署且要求符合信创数据库改造通常canal组件已经无法兼容信创库上述链路必须改造兼容通常改造量比较大且牵扯到各个业务线周期会很长。因上述消费是异步如果理论上本地化部署能接受秒级别延迟本文探索给出一种收益较高技术解决方案只写一个组件服务业务线代码不需要改动代码。以人大金仓库为例子。二、方案一捕捉业务表记录变更需要记录变更记录CREATE TABLE table_audit_log ( id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, //主键没有含义根据信创库情况设置 table_name VARCHAR(100) NOT NULL, // 监听变更的表名 action VARCHAR(10) NOT NULL, //insert、update两个值delete有没有物理删除看情况扩展定义 record_id VARCHAR(100) NOT NULL, // 监听记录的主键 old_data JSON, //变更前记录 new_data JSON, //变更后记录 deleted INT DEFAULT 0, // 被处理后删除 created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );二触发器捕获原业务表记录变更触发器捕获变更记录注意下面sql属于示例不同信创库创建trigger方式不一定一样。特别指出在本地事务中触发器和当前事务是在同一个事务中不存在记录变更后回滚了触发器已经记录了变更的情况。另外一个特别注意上面table_audit_log表是不能在上面做触发器的死循环问题。CREATE OR REPLACE TRIGGER trigger_表名_insert BEFORE INSERT ON 表名 FOR EACH ROW BEGIN INSERT INTO table_audit_log (table_name, action, record_id, new_data) VALUES (表名, insert, 新记录主键, json_build_object( 变更内容 ) ); END;三动态创建触发器如果给出一个表需要在上面创建两个触发器trigger_表名_insert和trigger_表名_update并且这两个创建不同的信创库语法不一样需要做成动态。以人大金仓为例CREATE OR REPLACE TRIGGER trigger_%s_insert BEFORE INSERT ON %s FOR EACH ROW BEGIN INSERT INTO table_audit_log (table_name, action, record_id, new_data) VALUES (%s, insert, :NEW.%s, json_build_object( %s ) ); END;为了创建上面trigger需要下面方法, 注意如果字段类型需要特别处理那就要特别处理比如下面的TIMESTAMP类型public String generateInsertTrigger(String tableName, String primaryKeyColumnName, ListColumnInfo columnInfoList) { if(CollectionUtils.isEmpty(columnInfoList)) { return null; } MapString, ColumnInfo columnInfoMap tableForTriggerMeta.getColumnMap(tableName); StringBuilder sbNew new StringBuilder(); for (ColumnInfo columnInfo : columnInfoList) { if(columnInfo.getColumnSize() TriggerDefine.COLUMN_MAX) { //跳过记录大字段 continue; } String columnName columnInfo.getColumnName(); int dataType columnInfoMap.get(columnName).getDataType(); if(Types.TIMESTAMP dataType) { String newValue to_char(:NEW. columnName , YYYY-MM-DD HH24:MI:SS.US); sbNew.append().append(columnName).append().append(, ).append(newValue).append(,); } else { sbNew.append().append(columnName).append().append(, :NEW.).append(columnName).append(,); } } String json sbNew.substring(0, sbNew.length()-1); return String.format(insertTemplate, tableName, tableName, tableName, primaryKeyColumnName, json); }另外上述代码有个细节点如果一张表某字段过大可以无需加入记录这是因为某些MQ比如rocketmq如果单个消息过大超过4M会无法发送另外业务上一般也用不到这么大的字段需要在业务侧处理如果需要建议单独反查大字段。四表名和表字段元数据获取可以通过jdbc根据表名获取对应元数据public TableMeta getTableColumns(String tableName) { TableMeta tableMeta new TableMeta(); ListColumnInfo columns new ArrayList(); tableMeta.setColumnInfoList(columns); jdbcTemplate.execute((Connection connection) - { DatabaseMetaData metaData connection.getMetaData(); String schemaName connection.getCatalog(); tableMeta.setSchemaName(schemaName); tableMeta.setDatabaseProductName(metaData.getDatabaseProductName()); ResultSet resultSet metaData.getPrimaryKeys(null, null, tableName); String primaryKeyName null; while (resultSet.next()) { primaryKeyName resultSet.getString(COLUMN_NAME); tableMeta.setPrimaryKeyColumnName(primaryKeyName); } log.info(primaryKeyName{}, primaryKeyName); // 获取表字段信息 try (ResultSet columnsResult metaData.getColumns( null, // catalog 数据库 null, // schema pattern tableName, // table name % // column name pattern )) { while (columnsResult.next()) { ColumnInfo column new ColumnInfo(); column.setTableName(columnsResult.getString(TABLE_NAME)); column.setColumnName(columnsResult.getString(COLUMN_NAME)); column.setDataType(columnsResult.getInt(DATA_TYPE)); column.setTypeName(columnsResult.getString(TYPE_NAME)); column.setColumnSize(columnsResult.getInt(COLUMN_SIZE)); column.setDecimalDigits(columnsResult.getInt(DECIMAL_DIGITS)); column.setIsNullable(columnsResult.getString(IS_NULLABLE)); column.setColumnDefault(columnsResult.getString(COLUMN_DEF)); column.setRemarks(columnsResult.getString(REMARKS)); column.setAutoIncrement(YES.equals(columnsResult.getString(IS_AUTOINCREMENT))); columns.add(column); } } return null; }); return tableMeta; }五哪些表需要做trigger技术框架需要一张表来记录哪些表需要做监听,如果10张表需要监听table_for_trigger会插入10条记录create table table_for_trigger ( id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, //主键没有含义根据信创库情况设置 table_name VARCHAR(100) NOT NULL, // 监听变更的表名 topic_name VARCHAR(10) NOT NULL, // 发送到MQ的topic trigger_insert_md5 VARCHAR(100), // 创建trigger后用于判断trigger是否有变化有变化需要重新创建 trigger_update_md5 VARCHAR(100), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );六发送到MQ通过定时任务发送到mq参考实现里面使用的是elasticjobrocketmq。三、参考实现https://github.com/zhangyl/canal-message-agent此实现其实可以很容易扩展集成做成框架以支持各种信创库但当前并不十分完善如果有同学愿意完善加入可以考虑做成开源框架。下面提几个待完善点table_for_trigger表目前没有做操作界面增加监听表某张监听表字段有变化需要重启服务才能感知可以做成在操作界面操作甚至自动感知跳过记录大字段需要配置当前写死