From d34bee33595ef975f0831fda89ac0643e89183fb Mon Sep 17 00:00:00 2001 From: mcy Date: Fri, 16 Nov 2018 18:07:32 +0800 Subject: [PATCH 1/3] =?UTF-8?q?adapter=20=E9=85=8D=E7=BD=AE=E4=BF=AE?= =?UTF-8?q?=E6=94=B9=E5=8A=A8=E6=80=81=E5=8A=A0=E8=BD=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/support/MappingConfigsLoader.java | 25 +++ .../canal/client/adapter/es/ESAdapter.java | 40 +++-- .../adapter/es/monitor/ESConfigMonitor.java | 151 ++++++++++++++++++ .../adapter/es/service/ESSyncService.java | 3 +- .../es/test/sync/LabelSyncJoinSub2Test.java | 12 +- .../es/test/sync/LabelSyncJoinSubTest.java | 12 +- .../es/test/sync/RoleSyncJoinOne2Test.java | 8 +- .../es/test/sync/RoleSyncJoinOneTest.java | 18 +-- .../es/test/sync/UserSyncJoinOneTest.java | 10 +- .../es/test/sync/UserSyncSingleTest.java | 12 +- .../client/adapter/hbase/HbaseAdapter.java | 46 ++++-- .../hbase/monitor/HbaseConfigMonitor.java | 129 +++++++++++++++ ...tor.java => ApplicationConfigMonitor.java} | 4 +- .../canal/client/adapter/rdb/RdbAdapter.java | 58 ++++--- .../adapter/rdb/monitor/RdbConfigMonitor.java | 141 ++++++++++++++++ 15 files changed, 581 insertions(+), 88 deletions(-) create mode 100644 client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/monitor/ESConfigMonitor.java create mode 100644 client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/monitor/HbaseConfigMonitor.java rename client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/{ApplicationRunningMonitor.java => ApplicationConfigMonitor.java} (97%) create mode 100644 client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/monitor/RdbConfigMonitor.java diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MappingConfigsLoader.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MappingConfigsLoader.java index c437731b..a7c8d7d6 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MappingConfigsLoader.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MappingConfigsLoader.java @@ -43,4 +43,29 @@ public class MappingConfigsLoader { return configContentMap; } + + public static String loadConfig(String name) { + // 先取本地文件,再取类路径 + File filePath = new File(".." + File.separator + Constant.CONF_DIR + File.separator + name); + if (!filePath.exists()) { + URL url = MappingConfigsLoader.class.getClassLoader().getResource(""); + if (url != null) { + filePath = new File(url.getPath() + name); + } + } + if (filePath.exists()) { + String fileName = filePath.getName(); + if (!fileName.endsWith(".yml")) { + return null; + } + try (InputStream in = new FileInputStream(filePath)) { + byte[] bytes = new byte[in.available()]; + in.read(bytes); + return new String(bytes, StandardCharsets.UTF_8); + } catch (IOException e) { + throw new RuntimeException("Read mapping config: " + filePath.getAbsolutePath() + " error. ", e); + } + } + return null; + } } diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java index 24cc9029..18cc0c93 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java @@ -1,15 +1,14 @@ package com.alibaba.otter.canal.client.adapter.es; import java.net.InetAddress; -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; +import java.util.*; +import java.util.concurrent.ConcurrentHashMap; import java.util.regex.Matcher; import java.util.regex.Pattern; import javax.sql.DataSource; +import com.alibaba.otter.canal.client.adapter.es.monitor.ESConfigMonitor; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.client.transport.TransportClient; import org.elasticsearch.common.settings.Settings; @@ -37,12 +36,14 @@ import com.alibaba.otter.canal.client.adapter.support.*; @SPI("es") public class ESAdapter implements OuterAdapter { - private Map esSyncConfig = new LinkedHashMap<>(); // 文件名对应配置 - private Map> dbTableEsSyncConfig = new LinkedHashMap<>(); // schema-table对应配置 + private Map esSyncConfig = new ConcurrentHashMap<>(); // 文件名对应配置 + private Map> dbTableEsSyncConfig = new ConcurrentHashMap<>(); // schema-table对应配置 - private TransportClient transportClient; + private TransportClient transportClient; - private ESSyncService esSyncService; + private ESSyncService esSyncService; + + private ESConfigMonitor esConfigMonitor; public TransportClient getTransportClient() { return transportClient; @@ -56,7 +57,7 @@ public class ESAdapter implements OuterAdapter { return esSyncConfig; } - public Map> getDbTableEsSyncConfig() { + public Map> getDbTableEsSyncConfig() { return dbTableEsSyncConfig; } @@ -73,7 +74,9 @@ public class ESAdapter implements OuterAdapter { } }); - for (ESSyncConfig config : esSyncConfig.values()) { + for (Map.Entry entry : esSyncConfig.entrySet()) { + String configName = entry.getKey(); + ESSyncConfig config = entry.getValue(); SchemaItem schemaItem = SqlParser.parse(config.getEsMapping().getSql()); config.getEsMapping().setSchemaItem(schemaItem); @@ -89,9 +92,9 @@ public class ESAdapter implements OuterAdapter { String schema = matcher.group(2); schemaItem.getAliasTableItems().values().forEach(tableItem -> { - List esSyncConfigs = dbTableEsSyncConfig - .computeIfAbsent(schema + "-" + tableItem.getTableName(), k -> new ArrayList<>()); - esSyncConfigs.add(config); + Map esSyncConfigMap = dbTableEsSyncConfig + .computeIfAbsent(schema + "-" + tableItem.getTableName(), k -> new HashMap<>()); + esSyncConfigMap.put(configName, config); }); } @@ -108,6 +111,9 @@ public class ESAdapter implements OuterAdapter { } ESTemplate esTemplate = new ESTemplate(transportClient); esSyncService = new ESSyncService(esTemplate); + + esConfigMonitor = new ESConfigMonitor(); + esConfigMonitor.init(this); } catch (Exception e) { throw new RuntimeException(e); } @@ -117,8 +123,8 @@ public class ESAdapter implements OuterAdapter { public void sync(Dml dml) { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = dbTableEsSyncConfig.get(database + "-" + table); - esSyncService.sync(esSyncConfigs, dml); + Map configMap = dbTableEsSyncConfig.get(database + "-" + table); + esSyncService.sync(configMap.values(), dml); } @Override @@ -192,6 +198,10 @@ public class ESAdapter implements OuterAdapter { @Override public String getDestination(String task) { + if (esConfigMonitor != null) { + esConfigMonitor.destroy(); + } + ESSyncConfig config = esSyncConfig.get(task); if (config != null) { return config.getDestination(); diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/monitor/ESConfigMonitor.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/monitor/ESConfigMonitor.java new file mode 100644 index 00000000..3b19b39e --- /dev/null +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/monitor/ESConfigMonitor.java @@ -0,0 +1,151 @@ +package com.alibaba.otter.canal.client.adapter.es.monitor; + +import java.io.File; +import java.util.HashMap; +import java.util.Map; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import org.apache.commons.io.filefilter.FileFilterUtils; +import org.apache.commons.io.monitor.FileAlterationListenerAdaptor; +import org.apache.commons.io.monitor.FileAlterationMonitor; +import org.apache.commons.io.monitor.FileAlterationObserver; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.yaml.snakeyaml.Yaml; + +import com.alibaba.druid.pool.DruidDataSource; +import com.alibaba.otter.canal.client.adapter.es.ESAdapter; +import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfig; +import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem; +import com.alibaba.otter.canal.client.adapter.es.config.SqlParser; +import com.alibaba.otter.canal.client.adapter.support.DatasourceConfig; +import com.alibaba.otter.canal.client.adapter.support.MappingConfigsLoader; +import com.alibaba.otter.canal.client.adapter.support.Util; + +public class ESConfigMonitor { + + private static final Logger logger = LoggerFactory.getLogger(ESConfigMonitor.class); + + private static final String adapterName = "es"; + + private ESAdapter esAdapter; + + private FileAlterationMonitor fileMonitor; + + public void init(ESAdapter esAdapter) { + this.esAdapter = esAdapter; + File confDir = Util.getConfDirPath(adapterName); + try { + FileAlterationObserver observer = new FileAlterationObserver(confDir, + FileFilterUtils.and(FileFilterUtils.fileFileFilter(), FileFilterUtils.suffixFileFilter("yml"))); + FileListener listener = new FileListener(); + observer.addListener(listener); + fileMonitor = new FileAlterationMonitor(3000, observer); + fileMonitor.start(); + + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + public void destroy() { + try { + fileMonitor.stop(); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private class FileListener extends FileAlterationListenerAdaptor { + + @Override + public void onFileCreate(File file) { + super.onFileCreate(file); + try { + // 加载新增的配置文件 + String configContent = MappingConfigsLoader.loadConfig(adapterName + File.separator + file.getName()); + ESSyncConfig config = new Yaml().loadAs(configContent, ESSyncConfig.class); + config.validate(); + addConfigToCache(file, config); + + logger.info("Add a new es mapping config: {} to canal adapter", file.getName()); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileChange(File file) { + super.onFileChange(file); + + try { + if (esAdapter.getEsSyncConfig().containsKey(file.getName())) { + // 加载配置文件 + String configContent = MappingConfigsLoader + .loadConfig(adapterName + File.separator + file.getName()); + ESSyncConfig config = new Yaml().loadAs(configContent, ESSyncConfig.class); + config.validate(); + if (esAdapter.getEsSyncConfig().containsKey(file.getName())) { + deleteConfigFromCache(file); + } + addConfigToCache(file, config); + + logger.info("Change a es mapping config: {} of canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileDelete(File file) { + super.onFileDelete(file); + + try { + if (esAdapter.getEsSyncConfig().containsKey(file.getName())) { + deleteConfigFromCache(file); + + logger.info("Delete a es mapping config: {} of canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private void addConfigToCache(File file, ESSyncConfig config) { + esAdapter.getEsSyncConfig().put(file.getName(), config); + + SchemaItem schemaItem = SqlParser.parse(config.getEsMapping().getSql()); + config.getEsMapping().setSchemaItem(schemaItem); + + DruidDataSource dataSource = DatasourceConfig.DATA_SOURCES.get(config.getDataSourceKey()); + if (dataSource == null || dataSource.getUrl() == null) { + throw new RuntimeException("No data source found: " + config.getDataSourceKey()); + } + Pattern pattern = Pattern.compile(".*:(.*)://.*/(.*)\\?.*$"); + Matcher matcher = pattern.matcher(dataSource.getUrl()); + if (!matcher.find()) { + throw new RuntimeException("Not found the schema of jdbc-url: " + config.getDataSourceKey()); + } + String schema = matcher.group(2); + + schemaItem.getAliasTableItems().values().forEach(tableItem -> { + Map esSyncConfigMap = esAdapter.getDbTableEsSyncConfig() + .computeIfAbsent(schema + "-" + tableItem.getTableName(), k -> new HashMap<>()); + esSyncConfigMap.put(file.getName(), config); + }); + + } + + private void deleteConfigFromCache(File file) { + esAdapter.getEsSyncConfig().remove(file.getName()); + for (Map configMap : esAdapter.getDbTableEsSyncConfig().values()) { + if (configMap != null) { + configMap.remove(file.getName()); + } + } + + } + } +} diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java index 73b1a6bc..397e3e5a 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.client.adapter.es.service; +import java.util.Collection; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -39,7 +40,7 @@ public class ESSyncService { this.esTemplate = esTemplate; } - public void sync(List esSyncConfigs, Dml dml) { + public void sync(Collection esSyncConfigs, Dml dml) { long begin = System.currentTimeMillis(); if (esSyncConfigs != null) { if (logger.isTraceEnabled()) { diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSub2Test.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSub2Test.java index 2d59f565..966b115b 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSub2Test.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSub2Test.java @@ -50,9 +50,9 @@ public class LabelSyncJoinSub2Test { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b;a_", response.getSource().get("_labels")); @@ -88,9 +88,9 @@ public class LabelSyncJoinSub2Test { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b;aa_", response.getSource().get("_labels")); @@ -120,9 +120,9 @@ public class LabelSyncJoinSub2Test { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b_", response.getSource().get("_labels")); diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSubTest.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSubTest.java index eafcfe75..e8cbbe59 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSubTest.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/LabelSyncJoinSubTest.java @@ -50,9 +50,9 @@ public class LabelSyncJoinSubTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b;a", response.getSource().get("_labels")); @@ -88,9 +88,9 @@ public class LabelSyncJoinSubTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b;aa", response.getSource().get("_labels")); @@ -120,9 +120,9 @@ public class LabelSyncJoinSubTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("b", response.getSource().get("_labels")); diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOne2Test.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOne2Test.java index c87669ba..0fba73d3 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOne2Test.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOne2Test.java @@ -48,9 +48,9 @@ public class RoleSyncJoinOne2Test { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("admin_", response.getSource().get("_role_name")); @@ -85,9 +85,9 @@ public class RoleSyncJoinOne2Test { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("admin3_", response.getSource().get("_role_name")); diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOneTest.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOneTest.java index a7a61d5a..9739a355 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOneTest.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinOneTest.java @@ -49,9 +49,9 @@ public class RoleSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("admin", response.getSource().get("_role_name")); @@ -86,9 +86,9 @@ public class RoleSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("admin2", response.getSource().get("_role_name")); @@ -124,9 +124,9 @@ public class RoleSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("operator", response.getSource().get("_role_name")); @@ -151,7 +151,7 @@ public class RoleSyncJoinOneTest { old2.put("role_id", 2L); dml2.setOld(oldList2); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml2); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml2); GetResponse response2 = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("admin2", response2.getSource().get("_role_name")); @@ -181,9 +181,9 @@ public class RoleSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertNull(response.getSource().get("_role_name")); diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncJoinOneTest.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncJoinOneTest.java index 532cbc89..c0e43b61 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncJoinOneTest.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncJoinOneTest.java @@ -20,7 +20,7 @@ public class UserSyncJoinOneTest { @Before public void init() { -// AdapterConfigs.put("es", "mytest_user_join_one.yml"); + // AdapterConfigs.put("es", "mytest_user_join_one.yml"); esAdapter = Common.init(); } @@ -50,9 +50,9 @@ public class UserSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("Eric_", response.getSource().get("_name")); @@ -86,9 +86,9 @@ public class UserSyncJoinOneTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("Eric2_", response.getSource().get("_name")); diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncSingleTest.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncSingleTest.java index bcc3b7f6..6bf4b75b 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncSingleTest.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/UserSyncSingleTest.java @@ -43,9 +43,9 @@ public class UserSyncSingleTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("Eric", response.getSource().get("_name")); @@ -76,9 +76,9 @@ public class UserSyncSingleTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertEquals("Eric2", response.getSource().get("_name")); @@ -106,9 +106,9 @@ public class UserSyncSingleTest { String database = dml.getDatabase(); String table = dml.getTable(); - List esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); + Map esSyncConfigs = esAdapter.getDbTableEsSyncConfig().get(database + "-" + table); - esAdapter.getEsSyncService().sync(esSyncConfigs, dml); + esAdapter.getEsSyncService().sync(esSyncConfigs.values(), dml); GetResponse response = esAdapter.getTransportClient().prepareGet("mytest_user", "_doc", "1").get(); Assert.assertNull(response.getSource()); diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/HbaseAdapter.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/HbaseAdapter.java index 9c0395f3..a5172e75 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/HbaseAdapter.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/HbaseAdapter.java @@ -22,6 +22,7 @@ import org.slf4j.LoggerFactory; import com.alibaba.otter.canal.client.adapter.OuterAdapter; import com.alibaba.otter.canal.client.adapter.hbase.config.MappingConfig; import com.alibaba.otter.canal.client.adapter.hbase.config.MappingConfigLoader; +import com.alibaba.otter.canal.client.adapter.hbase.monitor.HbaseConfigMonitor; import com.alibaba.otter.canal.client.adapter.hbase.service.HbaseEtlService; import com.alibaba.otter.canal.client.adapter.hbase.service.HbaseSyncService; import com.alibaba.otter.canal.client.adapter.hbase.support.HbaseTemplate; @@ -36,14 +37,24 @@ import com.alibaba.otter.canal.client.adapter.support.*; @SPI("hbase") public class HbaseAdapter implements OuterAdapter { - private static Logger logger = LoggerFactory.getLogger(HbaseAdapter.class); + private static Logger logger = LoggerFactory.getLogger(HbaseAdapter.class); - private Map hbaseMapping = new ConcurrentHashMap<>(); // 文件名对应配置 - private Map mappingConfigCache = new ConcurrentHashMap<>(); // 库名-表名对应配置 + private Map hbaseMapping = new ConcurrentHashMap<>(); // 文件名对应配置 + private Map> mappingConfigCache = new ConcurrentHashMap<>(); // 库名-表名对应配置 - private Connection conn; - private HbaseSyncService hbaseSyncService; - private HbaseTemplate hbaseTemplate; + private Connection conn; + private HbaseSyncService hbaseSyncService; + private HbaseTemplate hbaseTemplate; + + private HbaseConfigMonitor configMonitor; + + public Map getHbaseMapping() { + return hbaseMapping; + } + + public Map> getMappingConfigCache() { + return mappingConfigCache; + } @Override public void init(OuterAdapterConfig configuration) { @@ -57,11 +68,14 @@ public class HbaseAdapter implements OuterAdapter { hbaseMapping.put(key, mappingConfig); } }); - for (MappingConfig mappingConfig : hbaseMapping.values()) { - mappingConfigCache.put(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "." - + mappingConfig.getHbaseMapping().getDatabase() + "." - + mappingConfig.getHbaseMapping().getTable(), - mappingConfig); + for (Map.Entry entry : hbaseMapping.entrySet()) { + String configName = entry.getKey(); + MappingConfig mappingConfig = entry.getValue(); + String k = StringUtils.trimToEmpty(mappingConfig.getDestination()) + "." + + mappingConfig.getHbaseMapping().getDatabase() + "." + + mappingConfig.getHbaseMapping().getTable(); + Map configMap = mappingConfigCache.computeIfAbsent(k, k1 -> new HashMap<>()); + configMap.put(configName, mappingConfig); } Map properties = configuration.getProperties(); @@ -71,6 +85,9 @@ public class HbaseAdapter implements OuterAdapter { conn = ConnectionFactory.createConnection(hbaseConfig); hbaseTemplate = new HbaseTemplate(conn); hbaseSyncService = new HbaseSyncService(hbaseTemplate); + + configMonitor = new HbaseConfigMonitor(); + configMonitor.init(this); } catch (Exception e) { throw new RuntimeException(e); } @@ -84,8 +101,8 @@ public class HbaseAdapter implements OuterAdapter { String destination = StringUtils.trimToEmpty(dml.getDestination()); String database = dml.getDatabase(); String table = dml.getTable(); - MappingConfig config = mappingConfigCache.get(destination + "." + database + "." + table); - hbaseSyncService.sync(config, dml); + Map configMap = mappingConfigCache.get(destination + "." + database + "." + table); + configMap.values().forEach(config -> hbaseSyncService.sync(config, dml)); } @Override @@ -160,6 +177,9 @@ public class HbaseAdapter implements OuterAdapter { @Override public void destroy() { + if (configMonitor != null) { + configMonitor.destroy(); + } if (conn != null) { try { conn.close(); diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/monitor/HbaseConfigMonitor.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/monitor/HbaseConfigMonitor.java new file mode 100644 index 00000000..ae6ddef2 --- /dev/null +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/monitor/HbaseConfigMonitor.java @@ -0,0 +1,129 @@ +package com.alibaba.otter.canal.client.adapter.hbase.monitor; + +import java.io.File; +import java.util.HashMap; +import java.util.Map; + +import org.apache.commons.io.filefilter.FileFilterUtils; +import org.apache.commons.io.monitor.FileAlterationListenerAdaptor; +import org.apache.commons.io.monitor.FileAlterationMonitor; +import org.apache.commons.io.monitor.FileAlterationObserver; +import org.apache.commons.lang.StringUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.yaml.snakeyaml.Yaml; + +import com.alibaba.otter.canal.client.adapter.hbase.HbaseAdapter; +import com.alibaba.otter.canal.client.adapter.hbase.config.MappingConfig; +import com.alibaba.otter.canal.client.adapter.support.MappingConfigsLoader; +import com.alibaba.otter.canal.client.adapter.support.Util; + +public class HbaseConfigMonitor { + + private static final Logger logger = LoggerFactory.getLogger(HbaseConfigMonitor.class); + + private static final String adapterName = "hbase"; + + private HbaseAdapter hbaseAdapter; + + private FileAlterationMonitor fileMonitor; + + public void init(HbaseAdapter hbaseAdapter) { + this.hbaseAdapter = hbaseAdapter; + File confDir = Util.getConfDirPath(adapterName); + try { + FileAlterationObserver observer = new FileAlterationObserver(confDir, + FileFilterUtils.and(FileFilterUtils.fileFileFilter(), FileFilterUtils.suffixFileFilter("yml"))); + FileListener listener = new FileListener(); + observer.addListener(listener); + fileMonitor = new FileAlterationMonitor(3000, observer); + fileMonitor.start(); + + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + public void destroy() { + try { + fileMonitor.stop(); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private class FileListener extends FileAlterationListenerAdaptor { + + @Override + public void onFileCreate(File file) { + super.onFileCreate(file); + try { + // 加载新增的配置文件 + String configContent = MappingConfigsLoader.loadConfig(adapterName + File.separator + file.getName()); + MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class); + config.validate(); + addConfigToCache(file, config); + + logger.info("Add a new hbase mapping config: {} to canal adapter", file.getName()); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileChange(File file) { + super.onFileChange(file); + + try { + if (hbaseAdapter.getHbaseMapping().containsKey(file.getName())) { + // 加载配置文件 + String configContent = MappingConfigsLoader + .loadConfig(adapterName + File.separator + file.getName()); + MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class); + config.validate(); + if (hbaseAdapter.getHbaseMapping().containsKey(file.getName())) { + deleteConfigFromCache(file); + } + addConfigToCache(file, config); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileDelete(File file) { + super.onFileDelete(file); + + try { + if (hbaseAdapter.getHbaseMapping().containsKey(file.getName())) { + deleteConfigFromCache(file); + + logger.info("Delete a hbase mapping config: {} of canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private void addConfigToCache(File file, MappingConfig config) { + hbaseAdapter.getHbaseMapping().put(file.getName(), config); + Map configMap = hbaseAdapter.getMappingConfigCache() + .computeIfAbsent(StringUtils.trimToEmpty(config.getDestination()) + "." + + config.getHbaseMapping().getDatabase() + "." + config.getHbaseMapping().getTable(), + k1 -> new HashMap<>()); + configMap.put(file.getName(), config); + } + + private void deleteConfigFromCache(File file) { + + hbaseAdapter.getHbaseMapping().remove(file.getName()); + for (Map configMap : hbaseAdapter.getMappingConfigCache().values()) { + if (configMap != null) { + configMap.remove(file.getName()); + } + } + + } + } +} diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationRunningMonitor.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationConfigMonitor.java similarity index 97% rename from client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationRunningMonitor.java rename to client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationConfigMonitor.java index 0397464d..8d519184 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationRunningMonitor.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/monitor/ApplicationConfigMonitor.java @@ -22,9 +22,9 @@ import com.alibaba.otter.canal.adapter.launcher.loader.CanalAdapterService; import com.alibaba.otter.canal.client.adapter.support.Util; @Component -public class ApplicationRunningMonitor { +public class ApplicationConfigMonitor { - private static final Logger logger = LoggerFactory.getLogger(ApplicationRunningMonitor.class); + private static final Logger logger = LoggerFactory.getLogger(ApplicationConfigMonitor.class); @Resource private ContextRefresher contextRefresher; 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 4e7e2504..11fe5c2f 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 @@ -28,22 +28,31 @@ 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 HashMap<>(); // 文件名对应配置 + private Map> mappingConfigCache = new HashMap<>(); // 库名-表名对应配置 - private DruidDataSource dataSource; + private DruidDataSource dataSource; - private RdbSyncService rdbSyncService; + private RdbSyncService rdbSyncService; - private int commitSize = 3000; + private int commitSize = 3000; - private volatile boolean running = false; + private volatile boolean running = false; - private List dmlList = Collections.synchronizedList(new ArrayList<>()); - private Lock syncLock = new ReentrantLock(); - private ExecutorService executor = Executors.newFixedThreadPool(1); + private List dmlList = Collections + .synchronizedList(new ArrayList<>()); + private Lock syncLock = new ReentrantLock(); + private ExecutorService executor = Executors.newFixedThreadPool(1); + + public Map getRdbMapping() { + return rdbMapping; + } + + public Map> getMappingConfigCache() { + return mappingConfigCache; + } @Override public void init(OuterAdapterConfig configuration) { @@ -56,11 +65,15 @@ public class RdbAdapter implements OuterAdapter { rdbMapping.put(key, mappingConfig); } }); - for (MappingConfig mappingConfig : rdbMapping.values()) { - mappingConfigCache - .put(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "." - + mappingConfig.getDbMapping().getDatabase() + "." + mappingConfig.getDbMapping().getTable(), - mappingConfig); + 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); } Map properties = configuration.getProperties(); @@ -113,16 +126,19 @@ public class RdbAdapter implements OuterAdapter { String destination = StringUtils.trimToEmpty(dml.getDestination()); String database = dml.getDatabase(); String table = dml.getTable(); - MappingConfig config = mappingConfigCache.get(destination + "." + database + "." + table); + Map configMap = mappingConfigCache.get(destination + "." + database + "." + table); - List simpleDmlList = SimpleDml.dml2SimpleDml(dml, config); + if (configMap != null) { + configMap.values().forEach(config -> { + List simpleDmlList = SimpleDml.dml2SimpleDml(dml, config); - dmlList.addAll(simpleDmlList); + dmlList.addAll(simpleDmlList); - if (dmlList.size() > commitSize) { - sync(); + if (dmlList.size() > commitSize) { + sync(); + } + }); } - if (logger.isDebugEnabled()) { logger.debug("DML: {}", JSON.toJSONString(dml, SerializerFeature.WriteMapNullValue)); } diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/monitor/RdbConfigMonitor.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/monitor/RdbConfigMonitor.java new file mode 100644 index 00000000..d74cd884 --- /dev/null +++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/monitor/RdbConfigMonitor.java @@ -0,0 +1,141 @@ +package com.alibaba.otter.canal.client.adapter.rdb.monitor; + +import java.io.File; +import java.util.HashMap; +import java.util.Map; + +import org.apache.commons.io.filefilter.FileFilterUtils; +import org.apache.commons.io.monitor.FileAlterationListenerAdaptor; +import org.apache.commons.io.monitor.FileAlterationMonitor; +import org.apache.commons.io.monitor.FileAlterationObserver; +import org.apache.commons.lang.StringUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.yaml.snakeyaml.Yaml; + +import com.alibaba.otter.canal.client.adapter.rdb.RdbAdapter; +import com.alibaba.otter.canal.client.adapter.rdb.config.MappingConfig; +import com.alibaba.otter.canal.client.adapter.support.MappingConfigsLoader; +import com.alibaba.otter.canal.client.adapter.support.Util; + +public class RdbConfigMonitor { + + private static final Logger logger = LoggerFactory.getLogger(RdbConfigMonitor.class); + + private static final String adapterName = "rdb"; + + private String key; + + private RdbAdapter rdbAdapter; + + private FileAlterationMonitor fileMonitor; + + public void init(String key, RdbAdapter rdbAdapter) { + this.key = key; + this.rdbAdapter = rdbAdapter; + File confDir = Util.getConfDirPath(adapterName); + try { + FileAlterationObserver observer = new FileAlterationObserver(confDir, + FileFilterUtils.and(FileFilterUtils.fileFileFilter(), FileFilterUtils.suffixFileFilter("yml"))); + FileListener listener = new FileListener(); + observer.addListener(listener); + fileMonitor = new FileAlterationMonitor(3000, observer); + fileMonitor.start(); + + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + public void destroy() { + try { + fileMonitor.stop(); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private class FileListener extends FileAlterationListenerAdaptor { + + @Override + public void onFileCreate(File file) { + super.onFileCreate(file); + try { + // 加载新增的配置文件 + String configContent = MappingConfigsLoader.loadConfig(adapterName + File.separator + file.getName()); + MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class); + config.validate(); + if ((key == null && config.getOuterAdapterKey() == null) + || (key != null && key.equals(config.getOuterAdapterKey()))) { + addConfigToCache(file, config); + + logger.info("Add a new rdb mapping config: {} to canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileChange(File file) { + super.onFileChange(file); + + try { + if (rdbAdapter.getRdbMapping().containsKey(file.getName())) { + // 加载配置文件 + String configContent = MappingConfigsLoader.loadConfig(adapterName + File.separator + file.getName()); + MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class); + config.validate(); + if ((key == null && config.getOuterAdapterKey() == null) + || (key != null && key.equals(config.getOuterAdapterKey()))) { + if (rdbAdapter.getRdbMapping().containsKey(file.getName())) { + deleteConfigFromCache(file); + } + addConfigToCache(file, config); + } else { + // 不能修改outerAdapterKey + throw new RuntimeException("Outer adapter key not allowed modify"); + } + logger.info("Change a rdb mapping config: {} of canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + @Override + public void onFileDelete(File file) { + super.onFileDelete(file); + + try { + if (rdbAdapter.getRdbMapping().containsKey(file.getName())) { + deleteConfigFromCache(file); + + logger.info("Delete a rdb mapping config: {} of canal adapter", file.getName()); + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + + private void addConfigToCache(File file, MappingConfig config) { + rdbAdapter.getRdbMapping().put(file.getName(), config); + Map configMap = rdbAdapter.getMappingConfigCache() + .computeIfAbsent(StringUtils.trimToEmpty(config.getDestination()) + "." + + config.getDbMapping().getDatabase() + "." + config.getDbMapping().getTable(), + k1 -> new HashMap<>()); + configMap.put(file.getName(), config); + } + + private void deleteConfigFromCache(File file) { + + rdbAdapter.getRdbMapping().remove(file.getName()); + for (Map configMap : rdbAdapter.getMappingConfigCache().values()) { + if (configMap != null) { + configMap.remove(file.getName()); + } + } + + } + } +} From 991418ce6ef9f7edce57d85c691dc45b0c6695ed Mon Sep 17 00:00:00 2001 From: mcy Date: Fri, 16 Nov 2018 22:53:10 +0800 Subject: [PATCH 2/3] =?UTF-8?q?=E4=BF=AE=E6=94=B9rdb=E6=89=B9=E9=87=8F?= =?UTF-8?q?=E6=8F=90=E4=BA=A4=E5=88=86=E6=89=B9=E7=9A=84=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../launcher/config/RefresherConfig.java | 22 ---------------- .../canal/client/adapter/rdb/RdbAdapter.java | 25 ++++++++++++++----- .../adapter/rdb/service/RdbSyncService.java | 1 + 3 files changed, 20 insertions(+), 28 deletions(-) delete mode 100644 client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/RefresherConfig.java diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/RefresherConfig.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/RefresherConfig.java deleted file mode 100644 index 127b7ceb..00000000 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/RefresherConfig.java +++ /dev/null @@ -1,22 +0,0 @@ -package com.alibaba.otter.canal.adapter.launcher.config; - -import org.springframework.cloud.context.refresh.ContextRefresher; -import org.springframework.cloud.context.scope.refresh.RefreshScope; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; - -@Configuration -public class RefresherConfig { - - @Bean - public RefreshScope refreshScope() { - return new RefreshScope(); - } - - @Bean - public ContextRefresher contextRefresher(ConfigurableApplicationContext configurableApplicationContext, - RefreshScope refreshScope) { - return new ContextRefresher(configurableApplicationContext, refreshScope); - } -} 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 11fe5c2f..78001c03 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 @@ -5,6 +5,8 @@ import java.sql.SQLException; import java.util.*; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -20,6 +22,7 @@ import com.alibaba.fastjson.serializer.SerializerFeature; import com.alibaba.otter.canal.client.adapter.OuterAdapter; 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.RdbSyncService; import com.alibaba.otter.canal.client.adapter.rdb.support.SimpleDml; @@ -44,8 +47,11 @@ public class RdbAdapter implements OuterAdapter { private List dmlList = Collections .synchronizedList(new ArrayList<>()); private Lock syncLock = new ReentrantLock(); + private Condition condition = syncLock.newCondition(); private ExecutorService executor = Executors.newFixedThreadPool(1); + private RdbConfigMonitor rdbConfigMonitor; + public Map getRdbMapping() { return rdbMapping; } @@ -107,18 +113,21 @@ public class RdbAdapter implements OuterAdapter { executor.submit(() -> { while (running) { try { - int size1 = dmlList.size(); - Thread.sleep(3000); - int size2 = dmlList.size(); - if (size1 == size2) { + syncLock.lock(); + if (!condition.await(3, TimeUnit.SECONDS)) { // 超时提交 sync(); } } catch (Exception e) { logger.error(e.getMessage(), e); + } finally { + syncLock.unlock(); } } }); + + rdbConfigMonitor = new RdbConfigMonitor(); + rdbConfigMonitor.init(configuration.getKey(), this); } @Override @@ -131,10 +140,9 @@ public class RdbAdapter implements OuterAdapter { if (configMap != null) { configMap.values().forEach(config -> { List simpleDmlList = SimpleDml.dml2SimpleDml(dml, config); - dmlList.addAll(simpleDmlList); - if (dmlList.size() > commitSize) { + if (dmlList.size() >= commitSize) { sync(); } }); @@ -148,6 +156,7 @@ public class RdbAdapter implements OuterAdapter { try { syncLock.lock(); if (!dmlList.isEmpty()) { + condition.signal(); rdbSyncService.sync(dmlList); dmlList.clear(); } @@ -251,6 +260,10 @@ public class RdbAdapter implements OuterAdapter { @Override public void destroy() { running = false; + if (rdbConfigMonitor != null) { + rdbConfigMonitor.destroy(); + } + executor.shutdown(); if (rdbSyncService != null) { 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 9467b8b9..9f1f5469 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 @@ -10,6 +10,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicInteger; import javax.sql.DataSource; From fcdf95b5ecc3dc59e0d52f3de5fb594162cc6215 Mon Sep 17 00:00:00 2001 From: mcy Date: Sat, 17 Nov 2018 04:37:46 +0800 Subject: [PATCH 3/3] =?UTF-8?q?rdb=20=E6=89=B9=E9=87=8F=E6=8F=90=E4=BA=A4?= =?UTF-8?q?=E8=B6=85=E6=97=B6=E6=8F=90=E4=BA=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../canal/client/adapter/es/ESAdapter.java | 7 +++++-- .../adapter/es/service/ESSyncService.java | 1 - .../canal/client/adapter/rdb/RdbAdapter.java | 21 ++++++++----------- 3 files changed, 14 insertions(+), 15 deletions(-) diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java index 18cc0c93..35df8a76 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/ESAdapter.java @@ -1,14 +1,16 @@ package com.alibaba.otter.canal.client.adapter.es; import java.net.InetAddress; -import java.util.*; +import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.regex.Matcher; import java.util.regex.Pattern; import javax.sql.DataSource; -import com.alibaba.otter.canal.client.adapter.es.monitor.ESConfigMonitor; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.client.transport.TransportClient; import org.elasticsearch.common.settings.Settings; @@ -22,6 +24,7 @@ import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfig.ESMapping; import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfigLoader; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem; import com.alibaba.otter.canal.client.adapter.es.config.SqlParser; +import com.alibaba.otter.canal.client.adapter.es.monitor.ESConfigMonitor; import com.alibaba.otter.canal.client.adapter.es.service.ESEtlService; import com.alibaba.otter.canal.client.adapter.es.service.ESSyncService; import com.alibaba.otter.canal.client.adapter.es.support.ESTemplate; diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java index 397e3e5a..823a9d4c 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java @@ -14,7 +14,6 @@ import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.serializer.SerializerFeature; import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfig; import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfig.ESMapping; -import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfigLoader; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem.ColumnItem; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem.FieldItem; 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 78001c03..291ecbab 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 @@ -5,7 +5,6 @@ import java.sql.SQLException; import java.util.*; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -47,7 +46,6 @@ public class RdbAdapter implements OuterAdapter { private List dmlList = Collections .synchronizedList(new ArrayList<>()); private Lock syncLock = new ReentrantLock(); - private Condition condition = syncLock.newCondition(); private ExecutorService executor = Executors.newFixedThreadPool(1); private RdbConfigMonitor rdbConfigMonitor; @@ -112,16 +110,16 @@ public class RdbAdapter implements OuterAdapter { executor.submit(() -> { while (running) { + int beginSize = dmlList.size(); try { - syncLock.lock(); - if (!condition.await(3, TimeUnit.SECONDS)) { - // 超时提交 - sync(); - } - } catch (Exception e) { - logger.error(e.getMessage(), e); - } finally { - syncLock.unlock(); + Thread.sleep(3000); + } catch (InterruptedException e) { + // ignore + } + int endSize = dmlList.size(); + + if (endSize - beginSize < 300) { + sync(); } } }); @@ -156,7 +154,6 @@ public class RdbAdapter implements OuterAdapter { try { syncLock.lock(); if (!dmlList.isEmpty()) { - condition.signal(); rdbSyncService.sync(dmlList); dmlList.clear(); }