diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQDestination.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQDestination.java index e0d9124f..00d8b870 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQDestination.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQDestination.java @@ -7,6 +7,7 @@ package com.alibaba.otter.canal.connector.core.producer; * @version 1.0.0 */ public class MQDestination { + private String canalDestination; private String topic; private Integer partition; diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/spi/ExtensionLoader.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/spi/ExtensionLoader.java index e5c6708c..8985b4cb 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/spi/ExtensionLoader.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/spi/ExtensionLoader.java @@ -29,8 +29,7 @@ import org.slf4j.LoggerFactory; */ 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/"; @@ -38,8 +37,7 @@ 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<>(); @@ -180,8 +178,7 @@ 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); } } @@ -201,8 +198,7 @@ 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); } } @@ -268,8 +264,10 @@ public class ExtensionLoader { Map> extensionClasses = new HashMap<>(); if (spiDir != null && standbyDir != null) { - // 1. plugin folder,customized extension classLoader (jar_dir/plugin) - String dir = File.separator + this.getJarDirectoryPath() + spiDir; // + "plugin"; + // 1. plugin folder,customized extension classLoader + // (jar_dir/plugin) + String dir = File.separator + this.getJarDirectoryPath() + spiDir; // + + // "plugin"; File externalLibDir = new File(dir); if (!externalLibDir.exists()) { @@ -326,8 +324,7 @@ public class ExtensionLoader { try { BufferedReader reader = null; try { - reader = new BufferedReader( - new InputStreamReader(url.openStream(), StandardCharsets.UTF_8)); + reader = new BufferedReader(new InputStreamReader(url.openStream(), StandardCharsets.UTF_8)); String line = null; while ((line = reader.readLine()) != null) { final int ci = line.indexOf('#'); @@ -347,10 +344,12 @@ 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 { @@ -368,9 +367,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()); } @@ -380,9 +379,12 @@ 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); @@ -397,15 +399,13 @@ 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/deployer/pom.xml b/deployer/pom.xml index 47976554..bdeadc96 100644 --- a/deployer/pom.xml +++ b/deployer/pom.xml @@ -84,7 +84,24 @@ - + + org.apache.maven.plugins + maven-dependency-plugin + 2.10 + + + copy-dependencies-to-canal-deployer + package + + copy-dependencies + + + jar-with-dependencies + ${project.basedir}/target/canal/plugin + + + + org.apache.maven.plugins maven-assembly-plugin @@ -104,24 +121,6 @@ false - - org.apache.maven.plugins - maven-dependency-plugin - 2.10 - - - copy-dependencies-to-canal-client-service - package - - copy-dependencies - - - jar-with-dependencies - ${project.basedir}/target/canal/plugin - - - -