diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/AbstractCanalAdapterWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/AbstractCanalAdapterWorker.java
index db87bf01..fdf2a91e 100644
--- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/AbstractCanalAdapterWorker.java
+++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/AbstractCanalAdapterWorker.java
@@ -70,6 +70,7 @@ public abstract class AbstractCanalAdapterWorker {
});
return true;
} catch (Exception e) {
+ logger.error(e.getMessage(), e);
return false;
}
}));
@@ -108,6 +109,7 @@ public abstract class AbstractCanalAdapterWorker {
});
return true;
} catch (Exception e) {
+ logger.error(e.getMessage(), e);
return false;
}
}));
@@ -178,7 +180,7 @@ public abstract class AbstractCanalAdapterWorker {
/**
* 分批同步
- *
+ *
* @param dmls
* @param adapter
*/
diff --git a/client-adapter/launcher/src/main/resources/application.yml b/client-adapter/launcher/src/main/resources/application.yml
index 098bc1f3..4a5cac3b 100644
--- a/client-adapter/launcher/src/main/resources/application.yml
+++ b/client-adapter/launcher/src/main/resources/application.yml
@@ -15,7 +15,7 @@ spring:
canal.conf:
canalServerHost: 127.0.0.1:11111
# zookeeperHosts: slave1:2181
-# mqServers: slave1:6667 #or rocketmq
+# mqServers: 127.0.0.1:9092 #or rocketmq
# flatMessage: true
batchSize: 500
syncBatchSize: 1000
@@ -34,7 +34,7 @@ canal.conf:
groups:
- groupId: g1
outerAdapters:
- - name: logger
+# - name: logger
# - name: rdb
# key: oracle1
# properties:
@@ -42,15 +42,13 @@ canal.conf:
# jdbc.url: jdbc:oracle:thin:@localhost:49161:XE
# jdbc.username: mytest
# jdbc.password: m121212
-# - name: rdb
-# key: postgres1
-# properties:
-# jdbc.driverClassName: org.postgresql.Driver
-# jdbc.url: jdbc:postgresql://localhost:5432/postgres
-# jdbc.username: postgres
-# jdbc.password: 121212
-# threads: 1
-# commitSize: 3000
+ - name: rdb
+ key: mysql1
+ properties:
+ jdbc.driverClassName: com.mysql.jdbc.Driver
+ jdbc.url: jdbc:mysql://192.168.100.36/mytest?useUnicode=true
+ jdbc.username: root
+ jdbc.password: Ambari-123
# - name: hbase
# properties:
# hbase.zookeeper.quorum: 127.0.0.1
@@ -59,4 +57,4 @@ canal.conf:
# - name: es
# hosts: 127.0.0.1:9300
# properties:
-# cluster.name: elasticsearch
\ No newline at end of file
+# cluster.name: elasticsearch
diff --git a/client-adapter/rdb/pom.xml b/client-adapter/rdb/pom.xml
index 4ac00ffd..881b9399 100644
--- a/client-adapter/rdb/pom.xml
+++ b/client-adapter/rdb/pom.xml
@@ -24,6 +24,11 @@
1.19
provided
+
+ com.alibaba.fastsql
+ fastsql
+ 2.0.0_preview_644
+
mysql
@@ -100,4 +105,4 @@
-
\ No newline at end of file
+
diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/RdbAdapter.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/RdbAdapter.java
index ae38b124..cc2fd32e 100644
--- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/RdbAdapter.java
+++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/RdbAdapter.java
@@ -2,10 +2,10 @@ package com.alibaba.otter.canal.client.adapter.rdb;
import java.sql.Connection;
import java.sql.SQLException;
-import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -21,22 +21,26 @@ import com.alibaba.otter.canal.client.adapter.rdb.config.ConfigLoader;
import com.alibaba.otter.canal.client.adapter.rdb.config.MappingConfig;
import com.alibaba.otter.canal.client.adapter.rdb.monitor.RdbConfigMonitor;
import com.alibaba.otter.canal.client.adapter.rdb.service.RdbEtlService;
+import com.alibaba.otter.canal.client.adapter.rdb.service.RdbMirrorDbSyncService;
import com.alibaba.otter.canal.client.adapter.rdb.service.RdbSyncService;
+import com.alibaba.otter.canal.client.adapter.rdb.support.SyncUtil;
import com.alibaba.otter.canal.client.adapter.support.*;
@SPI("rdb")
public class RdbAdapter implements OuterAdapter {
- private static Logger logger = LoggerFactory.getLogger(RdbAdapter.class);
+ private static Logger logger = LoggerFactory.getLogger(RdbAdapter.class);
- private Map rdbMapping = new HashMap<>(); // 文件名对应配置
- private Map> mappingConfigCache = new HashMap<>(); // 库名-表名对应配置
+ private Map rdbMapping = new ConcurrentHashMap<>(); // 文件名对应配置
+ private Map> mappingConfigCache = new ConcurrentHashMap<>(); // 库名-表名对应配置
+ private Map mirrorDbConfigCache = new ConcurrentHashMap<>(); // 镜像库配置
private DruidDataSource dataSource;
private RdbSyncService rdbSyncService;
+ private RdbMirrorDbSyncService rdbMirrorDbSyncService;
- private ExecutorService executor = Executors.newFixedThreadPool(1);
+ private ExecutorService executor = Executors.newFixedThreadPool(1);
private RdbConfigMonitor rdbConfigMonitor;
@@ -62,12 +66,18 @@ public class RdbAdapter implements OuterAdapter {
for (Map.Entry entry : rdbMapping.entrySet()) {
String configName = entry.getKey();
MappingConfig mappingConfig = entry.getValue();
- Map configMap = mappingConfigCache
- .computeIfAbsent(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
- + mappingConfig.getDbMapping().getDatabase() + "."
- + mappingConfig.getDbMapping().getTable(),
- k1 -> new HashMap<>());
- configMap.put(configName, mappingConfig);
+ if (!mappingConfig.getDbMapping().isMirrorDb()) {
+ Map configMap = mappingConfigCache.computeIfAbsent(
+ StringUtils.trimToEmpty(mappingConfig.getDestination()) + "." + mappingConfig.getDbMapping()
+ .getDatabase() + "." + mappingConfig.getDbMapping().getTable(),
+ k1 -> new ConcurrentHashMap<>());
+ configMap.put(configName, mappingConfig);
+ } else {
+ // mirrorDB
+ mirrorDbConfigCache.put(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
+ + mappingConfig.getDbMapping().getDatabase(),
+ mappingConfig);
+ }
}
Map properties = configuration.getProperties();
@@ -96,6 +106,10 @@ public class RdbAdapter implements OuterAdapter {
dataSource,
threads != null ? Integer.valueOf(threads) : null);
+ rdbMirrorDbSyncService = new RdbMirrorDbSyncService(mirrorDbConfigCache,
+ dataSource,
+ threads != null ? Integer.valueOf(threads) : null);
+
rdbConfigMonitor = new RdbConfigMonitor();
rdbConfigMonitor.init(configuration.getKey(), this);
}
@@ -103,6 +117,7 @@ public class RdbAdapter implements OuterAdapter {
@Override
public void sync(List dmls) {
rdbSyncService.sync(dmls);
+ rdbMirrorDbSyncService.sync(dmls);
}
@Override
@@ -157,7 +172,7 @@ public class RdbAdapter implements OuterAdapter {
public Map count(String task) {
MappingConfig config = rdbMapping.get(task);
MappingConfig.DbMapping dbMapping = config.getDbMapping();
- String sql = "SELECT COUNT(1) AS cnt FROM " + dbMapping.getTargetTable();
+ String sql = "SELECT COUNT(1) AS cnt FROM " + SyncUtil.dbTable(dbMapping);
Connection conn = null;
Map res = new LinkedHashMap<>();
try {
@@ -183,7 +198,7 @@ public class RdbAdapter implements OuterAdapter {
}
}
}
- res.put("targetTable", dbMapping.getTargetTable());
+ res.put("targetTable", SyncUtil.dbTable(dbMapping));
return res;
}
diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/config/MappingConfig.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/config/MappingConfig.java
index 1d32147a..11594cde 100644
--- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/config/MappingConfig.java
+++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/config/MappingConfig.java
@@ -1,8 +1,6 @@
package com.alibaba.otter.canal.client.adapter.rdb.config;
-import java.util.LinkedHashSet;
import java.util.Map;
-import java.util.Set;
/**
* RDB表映射配置
@@ -66,30 +64,37 @@ public class MappingConfig {
if (dbMapping.database == null || dbMapping.database.isEmpty()) {
throw new NullPointerException("dbMapping.database");
}
- if (dbMapping.table == null || dbMapping.table.isEmpty()) {
+ if (!dbMapping.isMirrorDb() && (dbMapping.table == null || dbMapping.table.isEmpty())) {
throw new NullPointerException("dbMapping.table");
}
- if (dbMapping.targetTable == null || dbMapping.targetTable.isEmpty()) {
+ if (!dbMapping.isMirrorDb() && (dbMapping.targetTable == null || dbMapping.targetTable.isEmpty())) {
throw new NullPointerException("dbMapping.targetTable");
}
}
public static class DbMapping {
- private String database; // 数据库名或schema名
- private String table; // 表面名
- private Map targetPk; // 目标表主键字段
- private boolean mapAll = false; // 映射所有字段
- private String targetTable; // 目标表名
- private Map targetColumns; // 目标表字段映射
+ private Boolean mirrorDb = false; // 是否镜像库
+ private String database; // 数据库名或schema名
+ private String table; // 表名
+ private Map targetPk; // 目标表主键字段
+ private Boolean mapAll = false; // 映射所有字段
+ private String targetDb; // 目标库名
+ private String targetTable; // 目标表名
+ private Map targetColumns; // 目标表字段映射
- private String etlCondition; // etl条件sql
+ private String etlCondition; // etl条件sql
- private Set families = new LinkedHashSet<>(); // column family列表
private int readBatch = 5000;
- private int commitBatch = 5000; // etl等批量提交大小
+ private int commitBatch = 5000; // etl等批量提交大小
- // private volatile Map allColumns; // mapAll为true,自动设置改字段
+ public boolean isMirrorDb() {
+ return mirrorDb == null ? false : mirrorDb;
+ }
+
+ public void setMirrorDb(Boolean mirrorDb) {
+ this.mirrorDb = mirrorDb;
+ }
public String getDatabase() {
return database;
@@ -116,13 +121,21 @@ public class MappingConfig {
}
public boolean isMapAll() {
- return mapAll;
+ return mapAll == null ? false : mapAll;
}
- public void setMapAll(boolean mapAll) {
+ public void setMapAll(Boolean mapAll) {
this.mapAll = mapAll;
}
+ public String getTargetDb() {
+ return targetDb;
+ }
+
+ public void setTargetDb(String targetDb) {
+ this.targetDb = targetDb;
+ }
+
public String getTargetTable() {
return targetTable;
}
@@ -147,14 +160,6 @@ public class MappingConfig {
this.etlCondition = etlCondition;
}
- public Set getFamilies() {
- return families;
- }
-
- public void setFamilies(Set families) {
- this.families = families;
- }
-
public int getReadBatch() {
return readBatch;
}
@@ -170,13 +175,5 @@ public class MappingConfig {
public void setCommitBatch(int commitBatch) {
this.commitBatch = commitBatch;
}
-
- // public Map getAllColumns() {
- // return allColumns;
- // }
- //
- // public void setAllColumns(Map allColumns) {
- // this.allColumns = allColumns;
- // }
}
}
diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbEtlService.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbEtlService.java
index a42e5f3c..22f1bc03 100644
--- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbEtlService.java
+++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbEtlService.java
@@ -109,7 +109,7 @@ public class RdbEtlService {
logger.info(
dbMapping.getTable() + " etl completed in: " + (System.currentTimeMillis() - start) / 1000 + "s!");
- etlResult.setResultMessage("导入目标表 " + dbMapping.getTargetTable() + " 数据:" + successCount.get() + " 条");
+ etlResult.setResultMessage("导入目标表 " + SyncUtil.dbTable(dbMapping) + " 数据:" + successCount.get() + " 条");
} catch (Exception e) {
logger.error(e.getMessage(), e);
errMsg.add(hbaseTable + " etl failed! ==>" + e.getMessage());
@@ -187,7 +187,7 @@ public class RdbEtlService {
// }
StringBuilder insertSql = new StringBuilder();
- insertSql.append("INSERT INTO ").append(dbMapping.getTargetTable()).append(" (");
+ insertSql.append("INSERT INTO ").append(SyncUtil.dbTable(dbMapping)).append(" (");
columnsMap
.forEach((targetColumnName, srcColumnName) -> insertSql.append(targetColumnName).append(","));
@@ -209,7 +209,7 @@ public class RdbEtlService {
// 删除数据
Map values = new LinkedHashMap<>();
StringBuilder deleteSql = new StringBuilder(
- "DELETE FROM " + dbMapping.getTargetTable() + " WHERE ");
+ "DELETE FROM " + SyncUtil.dbTable(dbMapping) + " WHERE ");
appendCondition(dbMapping, deleteSql, values, rs);
try (PreparedStatement pstmt2 = connTarget.prepareStatement(deleteSql.toString())) {
int k = 1;
diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbMirrorDbSyncService.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbMirrorDbSyncService.java
new file mode 100644
index 00000000..7b3f9bdc
--- /dev/null
+++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbMirrorDbSyncService.java
@@ -0,0 +1,66 @@
+package com.alibaba.otter.canal.client.adapter.rdb.service;
+
+import java.io.StringWriter;
+import java.sql.Connection;
+import java.sql.Statement;
+import java.util.List;
+import java.util.Map;
+
+import javax.sql.DataSource;
+
+import com.alibaba.fastsql.sql.ast.SQLName;
+import com.alibaba.fastsql.sql.ast.SQLStatement;
+import com.alibaba.fastsql.sql.ast.statement.SQLExprTableSource;
+import com.alibaba.fastsql.sql.dialect.mysql.parser.MySqlStatementParser;
+import com.alibaba.fastsql.sql.dialect.mysql.visitor.MySqlOutputVisitor;
+import com.alibaba.fastsql.sql.dialect.mysql.visitor.MySqlSchemaStatVisitor;
+import com.alibaba.fastsql.sql.parser.SQLStatementParser;
+import org.apache.commons.lang.StringUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.alibaba.otter.canal.client.adapter.rdb.config.MappingConfig;
+import com.alibaba.otter.canal.client.adapter.support.Dml;
+
+public class RdbMirrorDbSyncService {
+
+ private static final Logger logger = LoggerFactory.getLogger(RdbMirrorDbSyncService.class);
+
+ private Map mirrorDbConfigCache; // 镜像库配置
+ private DataSource dataSource;
+
+ public RdbMirrorDbSyncService(Map mirrorDbConfigCache, DataSource dataSource,
+ Integer threads){
+ this.mirrorDbConfigCache = mirrorDbConfigCache;
+ this.dataSource = dataSource;
+ }
+
+ public void sync(List dmls) {
+ for (Dml dml : dmls) {
+ String destination = StringUtils.trimToEmpty(dml.getDestination());
+ String database = dml.getDatabase();
+ MappingConfig configMap = mirrorDbConfigCache.get(destination + "." + database);
+ if (configMap == null) {
+ continue;
+ }
+ if (dml.getSql() != null) {
+ // DDL
+ executeDdl(database, dml.getSql());
+ } else {
+ // DML
+ // TODO
+ }
+ }
+ }
+
+ private void executeDdl(String database, String sql) {
+ try (Connection conn = dataSource.getConnection(); Statement statement = conn.createStatement()) {
+ statement.execute(sql);
+ if (logger.isTraceEnabled()) {
+ logger.trace("Execute DDL sql: {} for database: {}", sql, database);
+ }
+ } catch (Exception e) {
+ logger.error(e.getMessage(), e);
+ }
+ }
+}
diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java
index ef80bb2f..38469ac4 100644
--- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java
+++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java
@@ -91,7 +91,9 @@ public class RdbSyncService {
dmlsPartition[hash].add(syncItem);
});
} else {
- int hash = Math.abs(Math.abs(config.getDbMapping().getTargetTable().hashCode()) % threads);
+ int hash = 0;
+ // Math.abs(Math.abs(config.getDbMapping().getTargetTable().hashCode()) %
+ // threads);
List singleDmls = SingleDml.dml2SingleDmls(dml);
singleDmls.forEach(singleDml -> {
SyncItem syncItem = new SyncItem(config, singleDml);
@@ -164,7 +166,7 @@ public class RdbSyncService {
Map columnsMap = SyncUtil.getColumnsMap(dbMapping, data);
StringBuilder insertSql = new StringBuilder();
- insertSql.append("INSERT INTO ").append(dbMapping.getTargetTable()).append(" (");
+ insertSql.append("INSERT INTO ").append(SyncUtil.dbTable(dbMapping)).append(" (");
columnsMap.forEach((targetColumnName, srcColumnName) -> insertSql.append(targetColumnName).append(","));
int len = insertSql.length();
@@ -228,7 +230,7 @@ public class RdbSyncService {
Map ctype = getTargetColumnType(batchExecutor.getConn(), config);
StringBuilder updateSql = new StringBuilder();
- updateSql.append("UPDATE ").append(dbMapping.getTargetTable()).append(" SET ");
+ updateSql.append("UPDATE ").append(SyncUtil.dbTable(dbMapping)).append(" SET ");
List