From d70dc8b88c5cd1a5506b191226fc89d3b8e72244 Mon Sep 17 00:00:00 2001 From: mcy Date: Thu, 25 Oct 2018 21:21:06 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A7=84=E8=8C=83=E6=B3=A8=E9=87=8A?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../canal/client/adapter/OuterAdapter.java | 2 +- .../adapter/support/AdapterConfigs.java | 17 ++++- .../adapter/support/CanalClientConfig.java | 41 +++++------- .../adapter/support/DatasourceConfig.java | 26 ++++--- .../canal/client/adapter/support/Dml.java | 20 +++--- .../client/adapter/support/EtlResult.java | 15 +++-- .../adapter/support/ExtensionLoader.java | 67 ++++++++----------- .../client/adapter/support/JdbcTypeUtil.java | 4 +- .../client/adapter/support/MessageUtil.java | 2 +- .../adapter/support/OuterAdapterConfig.java | 2 +- .../canal/client/adapter/support/Result.java | 8 ++- .../canal/client/adapter/support/SPI.java | 6 ++ .../adapter/hbase/config/MappingConfig.java | 45 ++++++------- .../hbase/config/MappingConfigLoader.java | 2 +- .../hbase/service/HbaseEtlService.java | 32 +++++++++ .../hbase/service/HbaseSyncService.java | 2 +- .../launcher/CanalAdapterApplication.java | 6 ++ .../adapter/launcher/common/EtlLock.java | 6 ++ .../adapter/launcher/common/SyncSwitch.java | 13 ++-- .../launcher/config/AdapterCanalConfig.java | 6 ++ .../launcher/config/AdapterConfig.java | 8 ++- .../launcher/config/CuratorClient.java | 6 ++ .../launcher/config/SpringContext.java | 6 ++ .../loader/AbstractCanalAdapterWorker.java | 2 +- .../loader/CanalAdapterKafkaWorker.java | 2 +- .../loader/CanalAdapterRocketMQWorker.java | 2 +- .../launcher/loader/CanalAdapterService.java | 6 ++ .../launcher/loader/CanalAdapterWorker.java | 2 +- .../adapter/launcher/rest/CommonRest.java | 6 ++ 29 files changed, 230 insertions(+), 132 deletions(-) diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/OuterAdapter.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/OuterAdapter.java index 286ba398..edb34049 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/OuterAdapter.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/OuterAdapter.java @@ -11,7 +11,7 @@ import com.alibaba.otter.canal.client.adapter.support.SPI; /** * 外部适配器接口 * - * @author machengyuan 2018-8-18 下午10:14:02 + * @author reweerma 2018-8-18 下午10:14:02 * @version 1.0.0 */ @SPI("logger") diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/AdapterConfigs.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/AdapterConfigs.java index 44d29ee9..0ccfb509 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/AdapterConfigs.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/AdapterConfigs.java @@ -1,13 +1,26 @@ package com.alibaba.otter.canal.client.adapter.support; +import java.util.HashMap; import java.util.LinkedHashSet; import java.util.Map; import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; +/** + * 适配器配置集合, 用于配置加载, 线程不安全 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ public class AdapterConfigs { - private static Map> configs = new ConcurrentHashMap<>(); + /** + * 类型下对应所有配置名, 如: + * hbase + * ┗━ mytest_person.yml + * ┗━ mytest_role.yml + * ┗━ mytest_department.yml + */ + private static Map> configs = new HashMap<>(); public static void put(String key, String value) { Set values = configs.get(key); diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java index 08246a42..29da946e 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java @@ -2,29 +2,26 @@ package com.alibaba.otter.canal.client.adapter.support; import java.util.ArrayList; import java.util.List; -import java.util.Properties; /** * 配置信息类 * - * @author machengyuan 2018-8-18 下午10:40:12 + * @author rewerma 2018-8-18 下午10:40:12 * @version 1.0.0 */ public class CanalClientConfig { - private String canalServerHost; + private String canalServerHost; // 单机模式下canal server的 ip:port - private String zookeeperHosts; + private String zookeeperHosts; // 集群模式下的zk地址, 如果配置了单机地址则以单机为准!! - private Properties properties; + private String bootstrapServers; // kafka or rocket mq 地址 - private String bootstrapServers; + private Boolean flatMessage = true; // 是否已flatMessage模式传输, 只适用于mq模式 - private List mqTopics; + private List mqTopics; // mq topic 列表 - private Boolean flatMessage = true; - - private List canalInstances; + private List canalInstances; // tcp 模式下 canal 实例列表, 与mq模式不能共存!! public String getCanalServerHost() { return canalServerHost; @@ -42,14 +39,6 @@ public class CanalClientConfig { this.zookeeperHosts = zookeeperHosts; } - public Properties getProperties() { - return properties; - } - - public void setProperties(Properties properties) { - this.properties = properties; - } - public String getBootstrapServers() { return bootstrapServers; } @@ -84,9 +73,9 @@ public class CanalClientConfig { public static class CanalInstance { - private String instance; + private String instance; // 实例名 - private List adapterGroups; + private List adapterGroups; // 适配器分组列表 public String getInstance() { return instance; @@ -109,7 +98,7 @@ public class CanalClientConfig { public static class AdapterGroup { - private List outAdapters; + private List outAdapters; // 适配器列表 public List getOutAdapters() { return outAdapters; @@ -122,11 +111,11 @@ public class CanalClientConfig { public static class MQTopic { - private String mqMode; + private String mqMode; // mq模式 kafka or rocketMQ - private String topic; + private String topic; // topic名 - private List groups = new ArrayList<>(); + private List groups = new ArrayList<>(); // 分组列表 public String getMqMode() { return mqMode; @@ -155,11 +144,11 @@ public class CanalClientConfig { public static class Group { - private String groupId; + private String groupId; // group id // private List adapters = new ArrayList<>(); - private List outAdapters; + private List outAdapters; // 适配器配置列表 public String getGroupId() { return groupId; diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/DatasourceConfig.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/DatasourceConfig.java index cf2ec471..acc07dd1 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/DatasourceConfig.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/DatasourceConfig.java @@ -1,21 +1,27 @@ package com.alibaba.otter.canal.client.adapter.support; -import com.alibaba.druid.pool.DruidDataSource; - import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import com.alibaba.druid.pool.DruidDataSource; + +/** + * 数据源配置 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ public class DatasourceConfig { - public final static Map DATA_SOURCES = new ConcurrentHashMap<>(); + public final static Map DATA_SOURCES = new ConcurrentHashMap<>(); // key对应的数据源 - private String driver = "com.mysql.jdbc.Driver"; - private String url; - private String database; - private String type = "mysql"; - private String username; - private String password; - private Integer maxActive = 3; + private String driver = "com.mysql.jdbc.Driver"; // 默认为mysql jdbc驱动 + private String url; // jdbc url + private String database; // jdbc database + private String type = "mysql"; // 类型, 默认为mysql + private String username; // jdbc username + private String password; // jdbc password + private Integer maxActive = 3; // 连接池最大连接数,默认为3 public String getDriver() { return driver; diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Dml.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Dml.java index 8e5af059..1d5357b4 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Dml.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Dml.java @@ -7,24 +7,24 @@ import java.util.Map; /** * DML操作转换对象 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewerma 2018-8-19 下午11:30:49 * @version 1.0.0 */ public class Dml implements Serializable { private static final long serialVersionUID = 2611556444074013268L; - private String destination; - private String database; - private String table; - private String type; + private String destination; // 对应canal的实例或者MQ的topic + private String database; // 数据库或schema + private String table; // 表名 + private String type; // 类型: INSERT UPDATE DELETE // binlog executeTime - private Long es; + private Long es; // 执行耗时 // dml build timeStamp - private Long ts; - private String sql; - private List> data; - private List> old; + private Long ts; // 同步时间 + private String sql; // 执行的sql, dml sql为空 + private List> data; // 数据列表 + private List> old; // 旧数据列表, 用于update, size和data的size一一对应 public String getDestination() { return destination; diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/EtlResult.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/EtlResult.java index 0004e51a..498af7a3 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/EtlResult.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/EtlResult.java @@ -2,14 +2,21 @@ package com.alibaba.otter.canal.client.adapter.support; import java.io.Serializable; +/** + * ETL的结果对象 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ public class EtlResult implements Serializable { + private static final long serialVersionUID = 4250522736289866505L; - private boolean succeeded = false; + private boolean succeeded = false; - private String resultMessage; + private String resultMessage; - private String errorMessage; + private String errorMessage; public boolean getSucceeded() { return succeeded; @@ -34,4 +41,4 @@ public class EtlResult implements Serializable { public void setErrorMessage(String errorMessage) { this.errorMessage = errorMessage; } -} \ No newline at end of file +} diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/ExtensionLoader.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/ExtensionLoader.java index b5de0d82..5b5c1bbb 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/ExtensionLoader.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/ExtensionLoader.java @@ -1,22 +1,11 @@ package com.alibaba.otter.canal.client.adapter.support; -import java.io.BufferedReader; -import java.io.File; -import java.io.FilenameFilter; -import java.io.IOException; -import java.io.InputStreamReader; +import java.io.*; import java.net.MalformedURLException; import java.net.URL; import java.net.URLClassLoader; import java.nio.file.Paths; -import java.util.Arrays; -import java.util.Collections; -import java.util.Enumeration; -import java.util.HashMap; -import java.util.Map; -import java.util.NoSuchElementException; -import java.util.Set; -import java.util.TreeSet; +import java.util.*; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.regex.Pattern; @@ -27,12 +16,13 @@ import org.slf4j.LoggerFactory; /** * SPI 类加载器 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewerma 2018-8-19 下午11:30:49 * @version 1.0.0 */ public class ExtensionLoader { - private static final Logger logger = LoggerFactory.getLogger(ExtensionLoader.class); + private static final Logger logger = LoggerFactory + .getLogger(ExtensionLoader.class); private static final String SERVICES_DIRECTORY = "META-INF/services/"; @@ -40,7 +30,8 @@ public class ExtensionLoader { private static final String DEFAULT_CLASSLOADER_POLICY = "internal"; - private static final Pattern NAME_SEPARATOR = Pattern.compile("\\s*[,]+\\s*"); + private static final Pattern NAME_SEPARATOR = Pattern + .compile("\\s*[,]+\\s*"); private static final ConcurrentMap, ExtensionLoader> EXTENSION_LOADERS = new ConcurrentHashMap<>(); @@ -271,7 +262,8 @@ public class ExtensionLoader { return instance; } catch (Throwable t) { throw new IllegalStateException("Extension instance(name: " + name + ", class: " + type - + ") could not be instantiated: " + t.getMessage(), t); + + ") could not be instantiated: " + t.getMessage(), + t); } } @@ -279,8 +271,8 @@ public class ExtensionLoader { if (type == null) throw new IllegalArgumentException("Extension type == null"); if (name == null) throw new IllegalArgumentException("Extension name == null"); Class clazz = getExtensionClasses().get(name); - if (clazz == null) throw new IllegalStateException("No such extension \"" + name + "\" for " + type.getName() - + "!"); + if (clazz == null) + throw new IllegalStateException("No such extension \"" + name + "\" for " + type.getName() + "!"); return clazz; } @@ -342,8 +334,8 @@ public class ExtensionLoader { logger.info("extension classpath dir: " + dir); File externalLibDir = new File(dir); if (!externalLibDir.exists()) { - externalLibDir = new File(File.separator + this.getJarDirectoryPath() + File.separator + "canal_client" - + File.separator + "lib"); + externalLibDir = new File( + File.separator + this.getJarDirectoryPath() + File.separator + "canal_client" + File.separator + "lib"); } if (externalLibDir.exists()) { File[] files = externalLibDir.listFiles(new FilenameFilter() { @@ -495,12 +487,10 @@ public class ExtensionLoader { // Class.forName(line, true, // classLoader); if (!type.isAssignableFrom(clazz)) { - throw new IllegalStateException("Error when load extension class(interface: " - + type - + ", class line: " - + clazz.getName() - + "), class " - + clazz.getName() + throw new IllegalStateException( + "Error when load extension class(interface: " + type + + ", class line: " + clazz.getName() + + "), class " + clazz.getName() + "is not subtype of interface."); } else { try { @@ -518,9 +508,9 @@ public class ExtensionLoader { extensionClasses.put(n, clazz); } else if (c != clazz) { cachedNames.remove(clazz); - throw new IllegalStateException("Duplicate extension " - + type.getName() - + " name " + n + " on " + throw new IllegalStateException( + "Duplicate extension " + type.getName() + " name " + + n + " on " + c.getName() + " and " + clazz.getName()); } @@ -530,12 +520,9 @@ public class ExtensionLoader { } } } catch (Throwable t) { - IllegalStateException e = new IllegalStateException("Failed to load extension class(interface: " - + type - + ", class line: " - + line - + ") in " - + url + IllegalStateException e = new IllegalStateException( + "Failed to load extension class(interface: " + type + ", class line: " + + line + ") in " + url + ", cause: " + t.getMessage(), t); @@ -550,13 +537,15 @@ public class ExtensionLoader { } } catch (Throwable t) { logger.error("Exception when load extension class(interface: " + type + ", class file: " + url - + ") in " + url, t); + + ") in " + url, + t); } } // end of while urls } } catch (Throwable t) { - logger.error("Exception when load extension class(interface: " + type + ", description file: " + fileName - + ").", t); + logger.error( + "Exception when load extension class(interface: " + type + ", description file: " + fileName + ").", + t); } } diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java index d4cb8a0e..c24ff7de 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java @@ -11,7 +11,7 @@ import org.slf4j.LoggerFactory; /** * 类型转换工具类 * - * @author machengyuan 2018-8-19 下午06:14:23 + * @author rewerma 2018-8-19 下午06:14:23 * @version 1.0.0 */ public class JdbcTypeUtil { @@ -30,7 +30,7 @@ public class JdbcTypeUtil { switch (jdbcType) { case Types.BIT: case Types.BOOLEAN: -// return Boolean.class; + // return Boolean.class; case Types.TINYINT: return Byte.TYPE; case Types.SMALLINT: diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java index 826ed917..ed5a4358 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java @@ -9,7 +9,7 @@ import com.alibaba.otter.canal.protocol.Message; /** * Message对象解析工具类 * - * @author machengyuan 2018-8-19 下午06:14:23 + * @author rewerma 2018-8-19 下午06:14:23 * @version 1.0.0 */ public class MessageUtil { diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/OuterAdapterConfig.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/OuterAdapterConfig.java index 3aed2f16..4e73791e 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/OuterAdapterConfig.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/OuterAdapterConfig.java @@ -5,7 +5,7 @@ import java.util.Map; /** * 外部适配器配置信息类 * - * @author machengyuan 2018-8-18 下午10:15:12 + * @author rewerma 2018-8-18 下午10:15:12 * @version 1.0.0 */ public class OuterAdapterConfig { diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Result.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Result.java index 03831f1f..a9602d05 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Result.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/Result.java @@ -3,9 +3,15 @@ package com.alibaba.otter.canal.client.adapter.support; import java.io.Serializable; import java.util.Date; +/** + * 用于rest的结果返回对象 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ public class Result implements Serializable { - public Integer code = 20000; + public Integer code = 20000; public Object data; public String message; public Date sysTime; diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/SPI.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/SPI.java index 7aa69938..b4f89b6a 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/SPI.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/SPI.java @@ -6,6 +6,12 @@ import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; +/** + * SPI装载器注解 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Documented @Retention(RetentionPolicy.RUNTIME) @Target({ ElementType.TYPE }) diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java index a5ea5707..8027fb51 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java @@ -1,21 +1,18 @@ package com.alibaba.otter.canal.client.adapter.hbase.config; -import java.util.LinkedHashMap; -import java.util.LinkedHashSet; -import java.util.Map; -import java.util.Objects; -import java.util.Set; +import java.util.*; /** * HBase表映射配置 * - * @author machengyuan 2018-8-21 下午06:45:49 + * @author rewerma 2018-8-21 下午06:45:49 * @version 1.0.0 */ public class MappingConfig { - private String dataSourceKey; - private HbaseOrm hbaseOrm; + private String dataSourceKey; // 数据源key + + private HbaseOrm hbaseOrm; // hbase映射配置 public String getDataSourceKey() { return dataSourceKey; @@ -130,7 +127,7 @@ public class MappingConfig { } public enum Mode { - STRING("STRING"), NATIVE("NATIVE"), PHOENIX("PHOENIX"); + STRING("STRING"), NATIVE("NATIVE"), PHOENIX("PHOENIX"); private String type; @@ -145,23 +142,23 @@ public class MappingConfig { public static class HbaseOrm { - private Mode mode = Mode.STRING; - private String destination; - private String database; - private String table; - private String hbaseTable; - private String family = "CF"; - private boolean uppercaseQualifier = true; - private boolean autoCreateTable = false; // 同步时HBase中表不存在的情况下自动建表 - private String rowKey; // 指定复合主键为rowKey - private Map columns; - private ColumnItem rowKeyColumn; - private String etlCondition; + private Mode mode = Mode.STRING; // hbase默认转换格式 + private String destination; // canal实例或MQ的topic + private String database; // 数据库名或schema名 + private String table; // 表面名 + private String hbaseTable; // hbase表名 + private String family = "CF"; // 默认统一column family + private boolean uppercaseQualifier = true; // 是否转大写 + private boolean autoCreateTable = false; // 同步时HBase中表不存在的情况下自动建表 + private String rowKey; // 指定复合主键为rowKey + private Map columns; // 字段映射 + private ColumnItem rowKeyColumn; // rowKey字段 + private String etlCondition; // etl条件sql - private Map columnItems = new LinkedHashMap<>(); - private Set families = new LinkedHashSet<>(); + private Map columnItems = new LinkedHashMap<>(); // 转换后的字段映射列表 + private Set families = new LinkedHashSet<>(); // column family列表 private int readBatch = 5000; - private int commitBatch = 5000; + private int commitBatch = 5000; // etl等批量提交大小 public Mode getMode() { return mode; diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfigLoader.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfigLoader.java index 0fb727d8..83ba56df 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfigLoader.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfigLoader.java @@ -18,7 +18,7 @@ import com.alibaba.otter.canal.client.adapter.support.AdapterConfigs; /** * HBase表映射配置加载器 * - * @author machengyuan 2018-8-21 下午06:45:49 + * @author rewerma 2018-8-21 下午06:45:49 * @version 1.0.0 */ public class MappingConfigLoader { diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java index ee110755..7e3c817a 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java @@ -24,6 +24,12 @@ import com.alibaba.otter.canal.client.adapter.support.EtlResult; import com.alibaba.otter.canal.client.adapter.support.JdbcTypeUtil; import com.google.common.base.Joiner; +/** + * HBase ETL 操作业务类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ public class HbaseEtlService { private static Logger logger = LoggerFactory.getLogger(HbaseEtlService.class); @@ -66,6 +72,12 @@ public class HbaseEtlService { } } + /** + * 建表 + * + * @param hbaseTemplate + * @param config + */ public static void createTable(HbaseTemplate hbaseTemplate, MappingConfig config) { try { // 判断hbase表是否存在,不存在则建表 @@ -79,6 +91,15 @@ public class HbaseEtlService { } } + /** + * 导入数据 + * + * @param ds 数据源 + * @param hbaseTemplate hbaseTemplate + * @param config 配置 + * @param params 筛选条件 + * @return 导入结果 + */ public static EtlResult importData(DataSource ds, HbaseTemplate hbaseTemplate, MappingConfig config, List params) { EtlResult etlResult = new EtlResult(); @@ -209,6 +230,17 @@ public class HbaseEtlService { return etlResult; } + /** + * 执行导入 + * + * @param ds + * @param sql + * @param hbaseOrm + * @param hbaseTemplate + * @param successCount + * @param errMsg + * @return + */ private static boolean executeSqlImport(DataSource ds, String sql, MappingConfig.HbaseOrm hbaseOrm, HbaseTemplate hbaseTemplate, AtomicLong successCount, List errMsg) { try { diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java index 1ff406fd..2edc396b 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java @@ -13,7 +13,7 @@ import com.alibaba.otter.canal.client.adapter.support.Dml; /** * HBase同步操作业务 * - * @author machengyuan 2018-8-21 下午06:45:49 + * @author rewerma 2018-8-21 下午06:45:49 * @version 1.0.0 */ public class HbaseSyncService { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/CanalAdapterApplication.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/CanalAdapterApplication.java index 56bbde8e..99140b0c 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/CanalAdapterApplication.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/CanalAdapterApplication.java @@ -3,6 +3,12 @@ package com.alibaba.otter.canal.adapter.launcher; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.builder.SpringApplicationBuilder; +/** + * 启动入口 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @SpringBootApplication public class CanalAdapterApplication { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/EtlLock.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/EtlLock.java index 67d0ac41..5fd5f3fe 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/EtlLock.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/EtlLock.java @@ -14,6 +14,12 @@ import org.springframework.stereotype.Component; import com.alibaba.otter.canal.adapter.launcher.config.CuratorClient; +/** + * Etl 同步锁 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component public class EtlLock { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/SyncSwitch.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/SyncSwitch.java index d96c753e..aa6e55bc 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/SyncSwitch.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/common/SyncSwitch.java @@ -19,16 +19,22 @@ import com.alibaba.otter.canal.adapter.launcher.config.AdapterCanalConfig; import com.alibaba.otter.canal.adapter.launcher.config.CuratorClient; import com.alibaba.otter.canal.common.utils.BooleanMutex; +/** + * 同步开关 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component public class SyncSwitch { private static final String SYN_SWITCH_ZK_NODE = "/sync-switch/"; - private static final Map LOCAL_LOCK = new ConcurrentHashMap<>(); + private static final Map LOCAL_LOCK = new ConcurrentHashMap<>(); - private static final Map DISTRIBUTED_LOCK = new ConcurrentHashMap<>(); + private static final Map DISTRIBUTED_LOCK = new ConcurrentHashMap<>(); - private static Mode mode = Mode.LOCAL; + private static Mode mode = Mode.LOCAL; @Resource private AdapterCanalConfig adapterCanalConfig; @@ -204,5 +210,4 @@ public class SyncSwitch { } } - } diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterCanalConfig.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterCanalConfig.java index 18a12d7c..2579e1f7 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterCanalConfig.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterCanalConfig.java @@ -9,6 +9,12 @@ import org.springframework.stereotype.Component; import com.alibaba.otter.canal.client.adapter.support.CanalClientConfig; +/** + * canal 的相关配置类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component @ConfigurationProperties(prefix = "canal.conf") public class AdapterCanalConfig extends CanalClientConfig { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterConfig.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterConfig.java index da37804f..f3e265f8 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterConfig.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/AdapterConfig.java @@ -13,6 +13,12 @@ import com.alibaba.druid.pool.DruidDataSource; import com.alibaba.otter.canal.client.adapter.support.AdapterConfigs; import com.alibaba.otter.canal.client.adapter.support.DatasourceConfig; +/** + * 适配器数据源及配置文件列表配置类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component @ConfigurationProperties(prefix = "adapter.conf") public class AdapterConfig { @@ -55,7 +61,7 @@ public class AdapterConfig { try { ds.init(); } catch (SQLException e) { - logger.error("#Failed to initial datasource: " + datasourceConfig.getUrl(), e); + logger.error("ERROR ## failed to initial datasource: " + datasourceConfig.getUrl(), e); } DatasourceConfig.DATA_SOURCES.put(entry.getKey(), ds); } diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java index b7944852..0ec60392 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java @@ -8,6 +8,12 @@ import org.apache.curator.framework.CuratorFrameworkFactory; import org.apache.curator.retry.ExponentialBackoffRetry; import org.springframework.stereotype.Component; +/** + * curator 配置类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component public class CuratorClient { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/SpringContext.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/SpringContext.java index ac860bea..987371e2 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/SpringContext.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/SpringContext.java @@ -5,6 +5,12 @@ import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.stereotype.Component; +/** + * spring util配置类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component public class SpringContext implements ApplicationContextAware { 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 bd3a47a4..cd818a9d 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 @@ -21,7 +21,7 @@ import com.alibaba.otter.canal.protocol.Message; /** * 适配器工作线程抽象类 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewerma 2018-8-19 下午11:30:49 * @version 1.0.0 */ public abstract class AbstractCanalAdapterWorker { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java index 22272e80..4c48125a 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java @@ -17,7 +17,7 @@ import com.alibaba.otter.canal.protocol.Message; /** * kafka对应的client适配器工作线程 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewerma 2018-8-19 下午11:30:49 * @version 1.0.0 */ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java index 09863a33..4632d7f0 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java @@ -16,7 +16,7 @@ import com.alibaba.otter.canal.protocol.Message; /** * kafka对应的client适配器工作线程 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewerma 2018-8-19 下午11:30:49 * @version 1.0.0 */ public class CanalAdapterRocketMQWorker extends AbstractCanalAdapterWorker { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterService.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterService.java index 03f515bf..2e535535 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterService.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterService.java @@ -15,6 +15,12 @@ import com.alibaba.otter.canal.adapter.launcher.config.AdapterCanalConfig; import com.alibaba.otter.canal.adapter.launcher.config.AdapterConfig; import com.alibaba.otter.canal.client.adapter.support.DatasourceConfig; +/** + * 适配器启动业务类 + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @Component public class CanalAdapterService { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterWorker.java index 1c324996..c728a15a 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterWorker.java @@ -15,7 +15,7 @@ import com.alibaba.otter.canal.protocol.Message; /** * 原生canal-server对应的client适配器工作线程 * - * @author machengyuan 2018-8-19 下午11:30:49 + * @author rewrema 2018-8-19 下午11:30:49 * @version 1.0.0 */ public class CanalAdapterWorker extends AbstractCanalAdapterWorker { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java index bee78e04..d2c163fb 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java @@ -17,6 +17,12 @@ import com.alibaba.otter.canal.client.adapter.support.EtlResult; import com.alibaba.otter.canal.client.adapter.support.ExtensionLoader; import com.alibaba.otter.canal.client.adapter.support.Result; +/** + * 适配器操作Rest + * + * @author rewerma @ 2018-10-20 + * @version 1.0.0 + */ @RestController public class CommonRest {