From 03a95221911cc94bcb7c8df8ad98ca9b642addb8 Mon Sep 17 00:00:00 2001 From: agapple Date: Tue, 13 Nov 2018 13:41:19 +0800 Subject: [PATCH 1/4] fixed event --- .../otter/canal/parse/inbound/mysql/MysqlEventParser.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java index e8ca45e3..3c06ee6b 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java @@ -651,7 +651,7 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE throw new CanalParseException("command : 'show master status' has an error! pls check. you need (at least one of) the SUPER,REPLICATION CLIENT privilege(s) for this operation"); } EntryPosition endPosition = new EntryPosition(fields.get(0), Long.valueOf(fields.get(1))); - if (isGTIDMode && fields.size() > 4) { + if (isGTIDMode() && fields.size() > 4) { endPosition.setGtid(fields.get(4)); } return endPosition; From 154ba422f2ae6eb5d819f06f5469e256b30f454c Mon Sep 17 00:00:00 2001 From: agapple Date: Tue, 13 Nov 2018 14:10:48 +0800 Subject: [PATCH 2/4] fixed plugin --- client-adapter/launcher/pom.xml | 2 +- .../launcher/src/main/assembly/dev.xml | 102 +++++++++--------- .../launcher/src/main/assembly/release.xml | 8 +- .../launcher/src/main/bin/startup.bat | 2 +- .../launcher/src/main/bin/startup.sh | 4 + .../inbound/mysql/dbsync/LogEventConvert.java | 7 +- .../inbound/mysql/dbsync/TableMetaCache.java | 12 +++ 7 files changed, 77 insertions(+), 60 deletions(-) diff --git a/client-adapter/launcher/pom.xml b/client-adapter/launcher/pom.xml index d48e6e71..df31b8c8 100644 --- a/client-adapter/launcher/pom.xml +++ b/client-adapter/launcher/pom.xml @@ -137,7 +137,7 @@ jar-with-dependencies - ${project.basedir}/target/canal-adapter/lib + ${project.basedir}/target/canal-adapter/plugin diff --git a/client-adapter/launcher/src/main/assembly/dev.xml b/client-adapter/launcher/src/main/assembly/dev.xml index 84962aec..1138f3f3 100644 --- a/client-adapter/launcher/src/main/assembly/dev.xml +++ b/client-adapter/launcher/src/main/assembly/dev.xml @@ -6,36 +6,36 @@ false - - . - / - - README* - - - - ./src/main/bin - bin - - **/* - - 0755 - - - ./src/main/resources - /conf - - **/* + + . + / + + README* + + + + ./src/main/bin + bin + + **/* + + 0755 + + + ./src/main/resources + /conf + + **/* - - - - ../elasticsearch/src/main/resources/es - /conf/es - - **/* - - + + + + ../elasticsearch/src/main/resources/es + /conf/es + + **/* + + ../hbase/src/main/resources/hbase /conf/hbase @@ -43,27 +43,31 @@ **/* - - ../rdb/src/main/resources/ - /conf + + ../rdb/src/main/resources/ + /conf META-INF/** - - - target - logs - - **/* - - - - - - lib - - junit:junit - - - + + + target + logs + + **/* + + + + ${project.basedir}/target/canal-adapter/plugin + /plugin/ + + + + + lib + + junit:junit + + + diff --git a/client-adapter/launcher/src/main/assembly/release.xml b/client-adapter/launcher/src/main/assembly/release.xml index bebad929..456b4a6a 100644 --- a/client-adapter/launcher/src/main/assembly/release.xml +++ b/client-adapter/launcher/src/main/assembly/release.xml @@ -58,13 +58,9 @@ **/* - - ${project.basedir}/target/canal-adapter/lib - /lib/ - - *-jar-with-dependencies.jar - + ${project.basedir}/target/canal-adapter/plugin + /plugin/ diff --git a/client-adapter/launcher/src/main/bin/startup.bat b/client-adapter/launcher/src/main/bin/startup.bat index 0c464dff..5c834a27 100755 --- a/client-adapter/launcher/src/main/bin/startup.bat +++ b/client-adapter/launcher/src/main/bin/startup.bat @@ -8,7 +8,7 @@ if "%OS%" == "Windows_NT" set ENV_PATH=%~dp0% set conf_dir=%ENV_PATH%\..\conf set CLASSPATH=%conf_dir% -set CLASSPATH=%conf_dir%\..\lib\*;%CLASSPATH% +set CLASSPATH=%conf_dir%\..\lib\*;%conf_dir%\..\plugin\*;%CLASSPATH% set JAVA_MEM_OPTS= -Xms128m -Xmx512m -XX:PermSize=128m set JAVA_OPTS_EXT= -Djava.awt.headless=true -Djava.net.preferIPv4Stack=true -Dapplication.codeset=UTF-8 -Dfile.encoding=UTF-8 diff --git a/client-adapter/launcher/src/main/bin/startup.sh b/client-adapter/launcher/src/main/bin/startup.sh index 69993483..538733e1 100644 --- a/client-adapter/launcher/src/main/bin/startup.sh +++ b/client-adapter/launcher/src/main/bin/startup.sh @@ -53,6 +53,10 @@ ADAPTER_OPTS="-DappName=canal-adapter" for i in $base/lib/*; do CLASSPATH=$i:"$CLASSPATH"; done + +for i in $base/plugin/*; + do CLASSPATH=$i:"$CLASSPATH"; +done CLASSPATH="$base/conf:$CLASSPATH"; echo "cd to $bin_abs_path for workaround relative path" diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java index 5c5432c6..1ba0930d 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java @@ -646,15 +646,16 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar continue; } - if (fieldMeta != null && existOptionalMetaData) { + if (fieldMeta != null && existOptionalMetaData && tableMetaCache.isOnTSDB()) { // check column info boolean check = StringUtils.equalsIgnoreCase(fieldMeta.getColumnName(), info.name); check &= (fieldMeta.isUnsigned() == info.unsigned); check &= (fieldMeta.isNullable() == info.nullable); if (!check) { - throw new CanalParseException("MySQL8.0 unmatch column metadata & pls submit issue , db : " - + fieldMeta.toString() + " , binlog : " + info.toString() + throw new CanalParseException("MySQL8.0 unmatch column metadata & pls submit issue , table : " + + tableMeta.getFullName() + ", db fieldMeta : " + + fieldMeta.toString() + " , binlog fieldMeta : " + info.toString() + " , on : " + event.getHeader().getLogFileName() + ":" + event.getHeader().getLogPos()); } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java index 4c1cef3c..69b8a3e9 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java @@ -39,6 +39,7 @@ public class TableMetaCache { public static final String EXTRA = "EXTRA"; private MysqlConnection connection; private boolean isOnRDS = false; + private boolean isOnTSDB = false; private TableMetaTSDB tableMetaTSDB; // 第一层tableId,第二层schema.table,解决tableId重复,对应多张表 @@ -67,6 +68,8 @@ public class TableMetaCache { } }); + } else { + isOnTSDB = true; } try { @@ -244,6 +247,15 @@ public class TableMetaCache { .toString(); } + + public boolean isOnTSDB() { + return isOnTSDB; + } + + public void setOnTSDB(boolean isOnTSDB) { + this.isOnTSDB = isOnTSDB; + } + public boolean isOnRDS() { return isOnRDS; } From 8c3ade2d50cc02b4e1e582a2518716ed3a726dce Mon Sep 17 00:00:00 2001 From: mcy Date: Tue, 13 Nov 2018 17:40:12 +0800 Subject: [PATCH 3/4] =?UTF-8?q?=E5=A4=96=E9=83=A8adapter=20plugin=E6=96=87?= =?UTF-8?q?=E4=BB=B6=E5=A4=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/support/ExtensionLoader.java | 96 +++---------------- .../support/URLClassExtensionLoader.java | 88 +++++++++++++++++ .../launcher/src/main/assembly/dev.xml | 4 - .../launcher/src/main/bin/startup.bat | 2 +- .../launcher/src/main/bin/startup.sh | 3 - 5 files changed, 100 insertions(+), 93 deletions(-) create mode 100644 client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/URLClassExtensionLoader.java 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 62ea78fd..b29bff4a 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 @@ -43,8 +43,6 @@ public class ExtensionLoader { private static final ConcurrentMap EXTENSION_KEY_INSTANCE = new ConcurrentHashMap<>(); - private static final ConcurrentMap> EXTENSION_KEY_INSTANCES = new ConcurrentHashMap<>(); - private final Class type; private final String classLoaderPolicy; @@ -180,6 +178,8 @@ public class ExtensionLoader { @SuppressWarnings("unchecked") private T createExtension(String name, String key) { + System.out.println("xxxxxxxxxxxxx"); + getExtensionClasses().forEach((k, v) -> logger.info("fffff: " + k + " " + v.getName())); Class clazz = getExtensionClasses().get(name); if (clazz == null) { throw new IllegalStateException("Extension instance(name: " + name + ", class: " + type @@ -210,6 +210,7 @@ public class ExtensionLoader { } } } + return classes; } @@ -255,13 +256,13 @@ public class ExtensionLoader { Map> extensionClasses = new HashMap>(); - // 1. lib folder,customized extension classLoader (jar_dir/lib) - String dir = File.separator + this.getJarDirectoryPath() + File.separator + "lib"; + // 1. plugin folder,customized extension classLoader (jar_dir/plugin) + String dir = File.separator + this.getJarDirectoryPath() + File.separator + "plugin"; File externalLibDir = new File(dir); if (!externalLibDir.exists()) { externalLibDir = new File(File.separator + this.getJarDirectoryPath() + File.separator + "canal-adapter" - + File.separator + "lib"); + + File.separator + "plugin"); } logger.info("extension classpath dir: " + externalLibDir.getAbsolutePath()); if (externalLibDir.exists()) { @@ -279,49 +280,7 @@ public class ExtensionLoader { URLClassLoader localClassLoader; if (classLoaderPolicy == null || "".equals(classLoaderPolicy) || DEFAULT_CLASSLOADER_POLICY.equalsIgnoreCase(classLoaderPolicy)) { - localClassLoader = new URLClassLoader(new URL[] { url }, parent) { - - @Override - public Class loadClass(String name) throws ClassNotFoundException { - Class c = findLoadedClass(name); - if (c != null) { - return c; - } - - if (name.startsWith("java.") || name.startsWith("org.slf4j.") - || name.startsWith("org.apache.logging") - || name.startsWith("org.apache.commons.logging.")) { - // || name.startsWith("org.apache.hadoop.")) - // { - c = super.loadClass(name); - } - if (c != null) return c; - - try { - // 先加载jar内的class,可避免jar冲突 - c = findClass(name); - } catch (ClassNotFoundException e) { - c = null; - } - if (c != null) { - return c; - } - - return super.loadClass(name); - } - - @Override - public Enumeration getResources(String name) throws IOException { - @SuppressWarnings("unchecked") - Enumeration[] tmp = (Enumeration[]) new Enumeration[2]; - - tmp[0] = findResources(name); // local class - // path first - // tmp[1] = super.getResources(name); - - return new CompoundEnumeration<>(tmp); - } - }; + localClassLoader = new URLClassExtensionLoader(new URL[] { url }); } else { localClassLoader = new URLClassLoader(new URL[] { url }, parent); } @@ -331,48 +290,15 @@ public class ExtensionLoader { } } } + // 只加载外部spi, 不加载classpath // 2. load inner extension class with default classLoader - ClassLoader classLoader = findClassLoader(); - loadFile(extensionClasses, CANAL_DIRECTORY, classLoader); - loadFile(extensionClasses, SERVICES_DIRECTORY, classLoader); + // ClassLoader classLoader = findClassLoader(); + // loadFile(extensionClasses, CANAL_DIRECTORY, classLoader); + // loadFile(extensionClasses, SERVICES_DIRECTORY, classLoader); return extensionClasses; } - public static class CompoundEnumeration implements Enumeration { - - private Enumeration[] enums; - private int index = 0; - - public CompoundEnumeration(Enumeration[] enums){ - this.enums = enums; - } - - private boolean next() { - while (this.index < this.enums.length) { - if (this.enums[this.index] != null && this.enums[this.index].hasMoreElements()) { - return true; - } - - ++this.index; - } - - return false; - } - - public boolean hasMoreElements() { - return this.next(); - } - - public E nextElement() { - if (!this.next()) { - throw new NoSuchElementException(); - } else { - return this.enums[this.index].nextElement(); - } - } - } - private void loadFile(Map> extensionClasses, String dir, ClassLoader classLoader) { String fileName = dir + type.getName(); try { diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/URLClassExtensionLoader.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/URLClassExtensionLoader.java new file mode 100644 index 00000000..725d23e8 --- /dev/null +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/URLClassExtensionLoader.java @@ -0,0 +1,88 @@ +package com.alibaba.otter.canal.client.adapter.support; + +import java.io.IOException; +import java.net.URL; +import java.net.URLClassLoader; +import java.util.Enumeration; +import java.util.NoSuchElementException; + +public class URLClassExtensionLoader extends URLClassLoader { + public URLClassExtensionLoader(URL[] urls) { + super(urls); + } + + @Override + public Class loadClass(String name) throws ClassNotFoundException { + Class c = findLoadedClass(name); + if (c != null) { + return c; + } + + if (name.startsWith("java.") || name.startsWith("org.slf4j.") + || name.startsWith("org.apache.logging") + || name.startsWith("org.apache.commons.logging.")) { + // || name.startsWith("org.apache.hadoop.")) + // { + c = super.loadClass(name); + } + if (c != null) return c; + + try { + // 先加载jar内的class,可避免jar冲突 + c = findClass(name); + } catch (ClassNotFoundException e) { + c = null; + } + if (c != null) { + return c; + } + + return super.loadClass(name); + } + + @Override + public Enumeration getResources(String name) throws IOException { + @SuppressWarnings("unchecked") + Enumeration[] tmp = (Enumeration[]) new Enumeration[2]; + + tmp[0] = findResources(name); // local class + // path first + // tmp[1] = super.getResources(name); + + return new CompoundEnumeration<>(tmp); + } + + private static class CompoundEnumeration implements Enumeration { + + private Enumeration[] enums; + private int index = 0; + + public CompoundEnumeration(Enumeration[] enums){ + this.enums = enums; + } + + private boolean next() { + while (this.index < this.enums.length) { + if (this.enums[this.index] != null && this.enums[this.index].hasMoreElements()) { + return true; + } + + ++this.index; + } + + return false; + } + + public boolean hasMoreElements() { + return this.next(); + } + + public E nextElement() { + if (!this.next()) { + throw new NoSuchElementException(); + } else { + return this.enums[this.index].nextElement(); + } + } + } +} diff --git a/client-adapter/launcher/src/main/assembly/dev.xml b/client-adapter/launcher/src/main/assembly/dev.xml index 1138f3f3..7beaf269 100644 --- a/client-adapter/launcher/src/main/assembly/dev.xml +++ b/client-adapter/launcher/src/main/assembly/dev.xml @@ -57,10 +57,6 @@ **/* - - ${project.basedir}/target/canal-adapter/plugin - /plugin/ - diff --git a/client-adapter/launcher/src/main/bin/startup.bat b/client-adapter/launcher/src/main/bin/startup.bat index 5c834a27..0c464dff 100755 --- a/client-adapter/launcher/src/main/bin/startup.bat +++ b/client-adapter/launcher/src/main/bin/startup.bat @@ -8,7 +8,7 @@ if "%OS%" == "Windows_NT" set ENV_PATH=%~dp0% set conf_dir=%ENV_PATH%\..\conf set CLASSPATH=%conf_dir% -set CLASSPATH=%conf_dir%\..\lib\*;%conf_dir%\..\plugin\*;%CLASSPATH% +set CLASSPATH=%conf_dir%\..\lib\*;%CLASSPATH% set JAVA_MEM_OPTS= -Xms128m -Xmx512m -XX:PermSize=128m set JAVA_OPTS_EXT= -Djava.awt.headless=true -Djava.net.preferIPv4Stack=true -Dapplication.codeset=UTF-8 -Dfile.encoding=UTF-8 diff --git a/client-adapter/launcher/src/main/bin/startup.sh b/client-adapter/launcher/src/main/bin/startup.sh index 538733e1..650fba7d 100644 --- a/client-adapter/launcher/src/main/bin/startup.sh +++ b/client-adapter/launcher/src/main/bin/startup.sh @@ -54,9 +54,6 @@ for i in $base/lib/*; do CLASSPATH=$i:"$CLASSPATH"; done -for i in $base/plugin/*; - do CLASSPATH=$i:"$CLASSPATH"; -done CLASSPATH="$base/conf:$CLASSPATH"; echo "cd to $bin_abs_path for workaround relative path" From 20062d861d2236ffa381eecd430ff789856c0e65 Mon Sep 17 00:00:00 2001 From: agapple Date: Tue, 13 Nov 2018 19:59:18 +0800 Subject: [PATCH 4/4] fixed test code --- .../canal/client/adapter/support/ExtensionLoader.java | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) 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 b29bff4a..860cb82b 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 @@ -2,13 +2,15 @@ package com.alibaba.otter.canal.client.adapter.support; import java.io.BufferedReader; import java.io.File; -import java.io.IOException; import java.io.InputStreamReader; import java.net.MalformedURLException; import java.net.URL; import java.net.URLClassLoader; import java.nio.file.Paths; -import java.util.*; +import java.util.Arrays; +import java.util.Enumeration; +import java.util.HashMap; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.regex.Pattern; @@ -178,8 +180,9 @@ public class ExtensionLoader { @SuppressWarnings("unchecked") private T createExtension(String name, String key) { - System.out.println("xxxxxxxxxxxxx"); - getExtensionClasses().forEach((k, v) -> logger.info("fffff: " + k + " " + v.getName())); + // System.out.println("xxxxxxxxxxxxx"); + // getExtensionClasses().forEach((k, v) -> logger.info("fffff: " + k + + // " " + v.getName())); Class clazz = getExtensionClasses().get(name); if (clazz == null) { throw new IllegalStateException("Extension instance(name: " + name + ", class: " + type