From e7f07d7d3e8cae9f1cd9ae6829a8677e2543b008 Mon Sep 17 00:00:00 2001 From: mcy Date: Thu, 23 Aug 2018 17:33:15 +0800 Subject: [PATCH 1/5] =?UTF-8?q?=E6=95=B4=E5=90=88canal.server=E5=92=8Ccana?= =?UTF-8?q?l.kafka?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- deployer/pom.xml | 1 + .../otter/canal/deployer/CanalConstants.java | 1 + .../otter/canal/deployer/CanalLauncher.java | 13 ++ deployer/src/main/resources/canal.properties | 2 + .../src/main/resources/kafka.yml | 6 +- kafka/pom.xml | 145 ------------------ kafka/src/main/assembly/dev.xml | 64 -------- kafka/src/main/assembly/release.xml | 64 -------- kafka/src/main/bin/startup.bat | 25 --- kafka/src/main/bin/startup.sh | 104 ------------- kafka/src/main/bin/stop.sh | 65 -------- .../otter/canal/kafka/CanalLauncher.java | 17 -- .../otter/canal/kafka/CanalServerStarter.java | 78 ---------- kafka/src/main/resources/logback.xml | 85 ---------- pom.xml | 1 - server/pom.xml | 23 +++ .../canal/kafka}/CanalKafkaProducer.java | 7 +- .../otter/canal/kafka}/CanalKafkaStarter.java | 31 ++-- .../otter/canal/kafka}/KafkaProperties.java | 2 +- .../otter/canal/kafka}/MessageSerializer.java | 14 +- .../canal/server/CanalServerStarter.java | 12 ++ 21 files changed, 81 insertions(+), 679 deletions(-) rename {kafka => deployer}/src/main/resources/kafka.yml (85%) delete mode 100644 kafka/pom.xml delete mode 100644 kafka/src/main/assembly/dev.xml delete mode 100644 kafka/src/main/assembly/release.xml delete mode 100755 kafka/src/main/bin/startup.bat delete mode 100644 kafka/src/main/bin/startup.sh delete mode 100644 kafka/src/main/bin/stop.sh delete mode 100644 kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalLauncher.java delete mode 100644 kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalServerStarter.java delete mode 100644 kafka/src/main/resources/logback.xml rename {kafka/src/main/java/com/alibaba/otter/canal/kafka/producer => server/src/main/java/com/alibaba/otter/canal/kafka}/CanalKafkaProducer.java (91%) rename {kafka/src/main/java/com/alibaba/otter/canal/kafka/producer => server/src/main/java/com/alibaba/otter/canal/kafka}/CanalKafkaStarter.java (83%) rename {kafka/src/main/java/com/alibaba/otter/canal/kafka/producer => server/src/main/java/com/alibaba/otter/canal/kafka}/KafkaProperties.java (98%) rename {kafka/src/main/java/com/alibaba/otter/canal/kafka/producer => server/src/main/java/com/alibaba/otter/canal/kafka}/MessageSerializer.java (83%) create mode 100644 server/src/main/java/com/alibaba/otter/canal/server/CanalServerStarter.java diff --git a/deployer/pom.xml b/deployer/pom.xml index 2de1ce1e..76777f03 100644 --- a/deployer/pom.xml +++ b/deployer/pom.xml @@ -40,6 +40,7 @@ **/canal.properties **/spring/** **/example/** + **/kafka.yml diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java index b328feb7..1e51990e 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java @@ -22,6 +22,7 @@ public class CanalConstants { public static final String CANAL_AUTO_SCAN = ROOT + "." + "auto.scan"; public static final String CANAL_AUTO_SCAN_INTERVAL = ROOT + "." + "auto.scan.interval"; public static final String CANAL_CONF_DIR = ROOT + "." + "conf.dir"; + public static final String CANAL_SERVER_MODE = ROOT + "." + "serverMode"; public static final String CANAL_DESTINATION_SPLIT = ","; public static final String GLOBAL_NAME = "global"; diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalLauncher.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalLauncher.java index 828f49e2..2524186a 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalLauncher.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalLauncher.java @@ -3,6 +3,8 @@ package com.alibaba.otter.canal.deployer; import java.io.FileInputStream; import java.util.Properties; +import com.alibaba.otter.canal.kafka.CanalKafkaStarter; +import com.alibaba.otter.canal.server.CanalServerStarter; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -51,6 +53,17 @@ public class CanalLauncher { } }); + + CanalServerStarter canalServerStarter = null; + String serverMode = properties.getProperty(CanalConstants.CANAL_SERVER_MODE, "tcp"); + if (serverMode.equalsIgnoreCase("kafka")) { + canalServerStarter = new CanalKafkaStarter(); + } else if (serverMode.equalsIgnoreCase("rocketMQ")) { + // 预留rocketMQ启动 + } + if (canalServerStarter != null) { + canalServerStarter.init(); + } } catch (Throwable e) { logger.error("## Something goes wrong when starting up the canal Server:", e); System.exit(0); diff --git a/deployer/src/main/resources/canal.properties b/deployer/src/main/resources/canal.properties index 9083477d..56c35564 100644 --- a/deployer/src/main/resources/canal.properties +++ b/deployer/src/main/resources/canal.properties @@ -8,6 +8,8 @@ canal.zkServers= # flush data to zk canal.zookeeper.flush.period = 1000 canal.withoutNetty = false +# tcp, kafka, rocketMQ +canal.serverMode = tcp # flush meta cursor/parse position to file canal.file.data.dir = ${canal.conf.dir} canal.file.flush.period = 1000 diff --git a/kafka/src/main/resources/kafka.yml b/deployer/src/main/resources/kafka.yml similarity index 85% rename from kafka/src/main/resources/kafka.yml rename to deployer/src/main/resources/kafka.yml index 2729fc27..93d4d557 100644 --- a/kafka/src/main/resources/kafka.yml +++ b/deployer/src/main/resources/kafka.yml @@ -12,8 +12,8 @@ canalDestinations: topic: example partition: # 一个destination可以对应多个topic -# topics: -# - topic: example -# partition: + #topics: + # - topic: example + # partition: diff --git a/kafka/pom.xml b/kafka/pom.xml deleted file mode 100644 index b2587c02..00000000 --- a/kafka/pom.xml +++ /dev/null @@ -1,145 +0,0 @@ - - - 4.0.0 - - canal - com.alibaba.otter - 1.1.0-SNAPSHOT - ../pom.xml - - com.alibaba.otter - canal.kafka - jar - canal kafka module for otter ${project.version} - - - com.alibaba.otter - canal.deployer - ${project.version} - - - - org.yaml - snakeyaml - 1.17 - - - org.apache.kafka - kafka_2.11 - 1.1.1 - - - org.slf4j - slf4j-log4j12 - - - - - - org.jboss.netty - netty - 3.2.2.Final - - - - - - - - maven-jar-plugin - - - true - - - **/logback.xml - **/canal.properties - **/spring/** - **/example/** - **/kafka.yml - - - - - - org.apache.maven.plugins - maven-assembly-plugin - - 2.2.1 - - - assemble - - single - - package - - - - false - false - - - - - - - - dev - - true - - env - !release - - - - - - - maven-assembly-plugin - - - - ${basedir}/src/main/assembly/dev.xml - - canal - ${project.build.directory} - - - - - - - - - release - - - env - release - - - - - - - maven-assembly-plugin - - - - ${basedir}/src/main/assembly/release.xml - - - ${project.artifactId}-${project.version} - - ${project.parent.build.directory} - - - - - - - diff --git a/kafka/src/main/assembly/dev.xml b/kafka/src/main/assembly/dev.xml deleted file mode 100644 index 2c3e3fff..00000000 --- a/kafka/src/main/assembly/dev.xml +++ /dev/null @@ -1,64 +0,0 @@ - - dist - - dir - - false - - - . - / - - README* - - - - ./src/main/bin - bin - - **/* - - 0755 - - - ../deployer/src/main/conf - /conf - - **/* - - - - ../deployer/src/main/resources - /conf - - **/* - - - logback.xml - - - - ./src/main/resources - /conf - - **/* - - - - target - logs - - **/* - - - - - - lib - - junit:junit - - - - diff --git a/kafka/src/main/assembly/release.xml b/kafka/src/main/assembly/release.xml deleted file mode 100644 index 0b103620..00000000 --- a/kafka/src/main/assembly/release.xml +++ /dev/null @@ -1,64 +0,0 @@ - - dist - - tar.gz - - false - - - . - / - - README* - - - - ./src/main/bin - bin - - **/* - - 0755 - - - ../deployer/src/main/conf - /conf - - **/* - - - - ../deployer/src/main/resources - /conf - - **/* - - - logback.xml - - - - ./src/main/resources - /conf - - **/* - - - - target - logs - - **/* - - - - - - lib - - junit:junit - - - - diff --git a/kafka/src/main/bin/startup.bat b/kafka/src/main/bin/startup.bat deleted file mode 100755 index 429cc653..00000000 --- a/kafka/src/main/bin/startup.bat +++ /dev/null @@ -1,25 +0,0 @@ -@echo off -@if not "%ECHO%" == "" echo %ECHO% -@if "%OS%" == "Windows_NT" setlocal - -set ENV_PATH=.\ -if "%OS%" == "Windows_NT" set ENV_PATH=%~dp0% - -set conf_dir=%ENV_PATH%\..\conf -set canal_conf=%conf_dir%\canal.properties -set logback_configurationFile=%conf_dir%\logback.xml - -set CLASSPATH=%conf_dir% -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 -set JAVA_DEBUG_OPT= -server -Xdebug -Xnoagent -Djava.compiler=NONE -Xrunjdwp:transport=dt_socket,address=9099,server=y,suspend=n -set CANAL_OPTS= -DappName=otter-canal -Dlogback.configurationFile="%logback_configurationFile%" -Dcanal.conf="%canal_conf%" - -set JAVA_OPTS= %JAVA_MEM_OPTS% %JAVA_OPTS_EXT% %JAVA_DEBUG_OPT% %CANAL_OPTS% - -set CMD_STR= java %JAVA_OPTS% -classpath "%CLASSPATH%" java %JAVA_OPTS% -classpath "%CLASSPATH%" com.alibaba.otter.canal.kafka.CanalLauncher -echo start cmd : %CMD_STR% - -java %JAVA_OPTS% -classpath "%CLASSPATH%" com.alibaba.otter.canal.kafka.CanalLauncher \ No newline at end of file diff --git a/kafka/src/main/bin/startup.sh b/kafka/src/main/bin/startup.sh deleted file mode 100644 index 184c4b3f..00000000 --- a/kafka/src/main/bin/startup.sh +++ /dev/null @@ -1,104 +0,0 @@ -#!/bin/bash - -current_path=`pwd` -case "`uname`" in - Linux) - bin_abs_path=$(readlink -f $(dirname $0)) - ;; - *) - bin_abs_path=`cd $(dirname $0); pwd` - ;; -esac -base=${bin_abs_path}/.. -canal_conf=$base/conf/canal.properties -logback_configurationFile=$base/conf/logback.xml -export LANG=en_US.UTF-8 -export BASE=$base - -if [ -f $base/bin/canal.pid ] ; then - echo "found canal.pid , Please run stop.sh first ,then startup.sh" 2>&2 - exit 1 -fi - -if [ ! -d $base/logs/canal ] ; then - mkdir -p $base/logs/canal -fi - -## set java path -if [ -z "$JAVA" ] ; then - JAVA=$(which java) -fi - -ALIBABA_JAVA="/usr/alibaba/java/bin/java" -TAOBAO_JAVA="/opt/taobao/java/bin/java" -if [ -z "$JAVA" ]; then - if [ -f $ALIBABA_JAVA ] ; then - JAVA=$ALIBABA_JAVA - elif [ -f $TAOBAO_JAVA ] ; then - JAVA=$TAOBAO_JAVA - else - echo "Cannot find a Java JDK. Please set either set JAVA or put java (>=1.5) in your PATH." 2>&2 - exit 1 - fi -fi - -case "$#" -in -0 ) - ;; -1 ) - var=$* - if [ -f $var ] ; then - canal_conf=$var - else - echo "THE PARAMETER IS NOT CORRECT.PLEASE CHECK AGAIN." - exit - fi;; -2 ) - var=$1 - if [ -f $var ] ; then - canal_conf=$var - else - if [ "$1" = "debug" ]; then - DEBUG_PORT=$2 - DEBUG_SUSPEND="n" - JAVA_DEBUG_OPT="-Xdebug -Xnoagent -Djava.compiler=NONE -Xrunjdwp:transport=dt_socket,address=$DEBUG_PORT,server=y,suspend=$DEBUG_SUSPEND" - fi - fi;; -* ) - echo "THE PARAMETERS MUST BE TWO OR LESS.PLEASE CHECK AGAIN." - exit;; -esac - -str=`file -L $JAVA | grep 64-bit` -if [ -n "$str" ]; then - JAVA_OPTS="-server -Xms2048m -Xmx3072m -Xmn1024m -XX:SurvivorRatio=2 -XX:PermSize=96m -XX:MaxPermSize=256m -Xss256k -XX:-UseAdaptiveSizePolicy -XX:MaxTenuringThreshold=15 -XX:+DisableExplicitGC -XX:+UseConcMarkSweepGC -XX:+CMSParallelRemarkEnabled -XX:+UseCMSCompactAtFullCollection -XX:+UseFastAccessorMethods -XX:+UseCMSInitiatingOccupancyOnly -XX:+HeapDumpOnOutOfMemoryError" -else - JAVA_OPTS="-server -Xms1024m -Xmx1024m -XX:NewSize=256m -XX:MaxNewSize=256m -XX:MaxPermSize=128m " -fi - -JAVA_OPTS=" $JAVA_OPTS -Djava.awt.headless=true -Djava.net.preferIPv4Stack=true -Dfile.encoding=UTF-8" -CANAL_OPTS="-DappName=otter-canal -Dlogback.configurationFile=$logback_configurationFile -Dcanal.conf=$canal_conf" - -if [ -e $canal_conf -a -e $logback_configurationFile ] -then - - for i in $base/lib/*; - do CLASSPATH=$i:"$CLASSPATH"; - done - CLASSPATH="$base/conf:$CLASSPATH"; - - echo "cd to $bin_abs_path for workaround relative path" - cd $bin_abs_path - - echo LOG CONFIGURATION : $logback_configurationFile - echo canal conf : $canal_conf - echo CLASSPATH :$CLASSPATH - $JAVA $JAVA_OPTS $JAVA_DEBUG_OPT $CANAL_OPTS -classpath .:$CLASSPATH com.alibaba.otter.canal.kafka.CanalLauncher 1>>$base/logs/canal/canal.log 2>&1 & - echo $! > $base/bin/canal.pid - - echo "cd to $current_path for continue" - cd $current_path -else - echo "canal conf("$canal_conf") OR log configration file($logback_configurationFile) is not exist,please create then first!" -fi diff --git a/kafka/src/main/bin/stop.sh b/kafka/src/main/bin/stop.sh deleted file mode 100644 index f398749c..00000000 --- a/kafka/src/main/bin/stop.sh +++ /dev/null @@ -1,65 +0,0 @@ -#!/bin/bash - -cygwin=false; -linux=false; -case "`uname`" in - CYGWIN*) - cygwin=true - ;; - Linux*) - linux=true - ;; -esac - -get_pid() { - STR=$1 - PID=$2 - if $cygwin; then - JAVA_CMD="$JAVA_HOME\bin\java" - JAVA_CMD=`cygpath --path --unix $JAVA_CMD` - JAVA_PID=`ps |grep $JAVA_CMD |awk '{print $1}'` - else - if $linux; then - if [ ! -z "$PID" ]; then - JAVA_PID=`ps -C java -f --width 1000|grep "$STR"|grep "$PID"|grep -v grep|awk '{print $2}'` - else - JAVA_PID=`ps -C java -f --width 1000|grep "$STR"|grep -v grep|awk '{print $2}'` - fi - else - if [ ! -z "$PID" ]; then - JAVA_PID=`ps aux |grep "$STR"|grep "$PID"|grep -v grep|awk '{print $2}'` - else - JAVA_PID=`ps aux |grep "$STR"|grep -v grep|awk '{print $2}'` - fi - fi - fi - echo $JAVA_PID; -} - -base=`dirname $0`/.. -pidfile=$base/bin/canal.pid -if [ ! -f "$pidfile" ];then - echo "canal is not running. exists" - exit -fi - -pid=`cat $pidfile` -if [ "$pid" == "" ] ; then - pid=`get_pid "appName=otter-canal"` -fi - -echo -e "`hostname`: stopping canal $pid ... " -kill $pid - -LOOPS=0 -while (true); -do - gpid=`get_pid "appName=otter-canal" "$pid"` - if [ "$gpid" == "" ] ; then - echo "Oook! cost:$LOOPS" - `rm $pidfile` - break; - fi - let LOOPS=LOOPS+1 - sleep 1 -done \ No newline at end of file diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalLauncher.java b/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalLauncher.java deleted file mode 100644 index 3212733c..00000000 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalLauncher.java +++ /dev/null @@ -1,17 +0,0 @@ -package com.alibaba.otter.canal.kafka; - -import com.alibaba.otter.canal.kafka.producer.CanalKafkaStarter; - -/** - * canal-kafka独立版本启动的入口类 - * - * @author machengyuan 2018-6-11 下午05:30:49 - * @version 1.0.0 - */ -public class CanalLauncher { - - public static void main(String[] args) { - CanalServerStarter.init(); - CanalKafkaStarter.init(); - } -} diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalServerStarter.java b/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalServerStarter.java deleted file mode 100644 index 573f0eac..00000000 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/CanalServerStarter.java +++ /dev/null @@ -1,78 +0,0 @@ -package com.alibaba.otter.canal.kafka; - -import java.io.FileInputStream; -import java.util.Properties; - -import org.apache.commons.lang.StringUtils; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.alibaba.otter.canal.deployer.CanalController; - -/** - * canal server 启动类 - * - * @author machengyuan 2018-6-11 下午05:30:49 - * @version 1.0.0 - */ -public class CanalServerStarter { - - private static final String CLASSPATH_URL_PREFIX = "classpath:"; - private static final Logger logger = LoggerFactory.getLogger(CanalServerStarter.class); - private volatile static boolean running = false; - - public static void init() { - try { - logger.info("## set default uncaught exception handler"); - setGlobalUncaughtExceptionHandler(); - - logger.info("## load canal configurations"); - String conf = System.getProperty("canal.conf", "classpath:canal.properties"); - Properties properties = new Properties(); - if (conf.startsWith(CLASSPATH_URL_PREFIX)) { - conf = StringUtils.substringAfter(conf, CLASSPATH_URL_PREFIX); - properties.load(CanalLauncher.class.getClassLoader().getResourceAsStream(conf)); - } else { - properties.load(new FileInputStream(conf)); - } - - logger.info("## start the canal server."); - final CanalController controller = new CanalController(properties); - controller.start(); - running = true; - logger.info("## the canal server is running now ......"); - Runtime.getRuntime().addShutdownHook(new Thread() { - - public void run() { - try { - logger.info("## stop the canal server"); - running = false; - controller.stop(); - } catch (Throwable e) { - logger.warn("##something goes wrong when stopping canal Server:", e); - } finally { - logger.info("## canal server is down."); - } - } - - }); - } catch (Throwable e) { - logger.error("## Something goes wrong when starting up the canal Server:", e); - System.exit(0); - } - } - - public static boolean isRunning() { - return running; - } - - private static void setGlobalUncaughtExceptionHandler() { - Thread.setDefaultUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() { - - @Override - public void uncaughtException(Thread t, Throwable e) { - logger.error("UnCaughtException", e); - } - }); - } -} diff --git a/kafka/src/main/resources/logback.xml b/kafka/src/main/resources/logback.xml deleted file mode 100644 index 34e7b1e5..00000000 --- a/kafka/src/main/resources/logback.xml +++ /dev/null @@ -1,85 +0,0 @@ - - - - - %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{56} - %msg%n - - - - - - - destination - canal - - - - ../logs/${destination}/${destination}.log - - - ../logs/${destination}/%d{yyyy-MM-dd}/${destination}-%d{yyyy-MM-dd}-%i.log.gz - - - 512MB - - 60 - - - - %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{56} - %msg%n - - - - - - - - - destination - canal - - - - ../logs/${destination}/meta.log - - - ../logs/${destination}/%d{yyyy-MM-dd}/meta-%d{yyyy-MM-dd}-%i.log.gz - - - 32MB - - 60 - - - - %d{yyyy-MM-dd HH:mm:ss.SSS} - %msg%n - - - - - - - - - - - - - - - - - - - - - - - - - - - - \ No newline at end of file diff --git a/pom.xml b/pom.xml index e7bb158e..64588a28 100644 --- a/pom.xml +++ b/pom.xml @@ -117,7 +117,6 @@ client deployer example - kafka kafka-client prometheus client-adapter diff --git a/server/pom.xml b/server/pom.xml index b7f83d56..3b8e3ddb 100644 --- a/server/pom.xml +++ b/server/pom.xml @@ -25,6 +25,29 @@ canal.instance.manager ${project.version} + + + org.yaml + snakeyaml + 1.17 + + + org.apache.kafka + kafka_2.11 + 1.1.1 + + + org.slf4j + slf4j-log4j12 + + + + + + org.jboss.netty + netty + 3.2.2.Final + diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaProducer.java b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java similarity index 91% rename from kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaProducer.java rename to server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java index 0bb15308..b6dba051 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaProducer.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.producer; +package com.alibaba.otter.canal.kafka; import java.io.IOException; import java.util.Properties; @@ -10,7 +10,6 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.alibaba.otter.canal.kafka.producer.KafkaProperties.Topic; import com.alibaba.otter.canal.protocol.Message; /** @@ -21,7 +20,7 @@ import com.alibaba.otter.canal.protocol.Message; */ public class CanalKafkaProducer { - private static final Logger logger = LoggerFactory.getLogger(CanalKafkaProducer.class); + private static final Logger logger = LoggerFactory.getLogger(CanalKafkaProducer.class); private Producer producer; @@ -49,7 +48,7 @@ public class CanalKafkaProducer { } } - public void send(Topic topic, Message message) throws IOException { + public void send(KafkaProperties.Topic topic, Message message) throws IOException { // set canal.instance.filter.transaction.entry = true // boolean valid = false; diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaStarter.java similarity index 83% rename from kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java rename to server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaStarter.java index 0eedfa7a..d244ce31 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaStarter.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.producer; +package com.alibaba.otter.canal.kafka; import java.io.FileInputStream; import java.util.List; @@ -10,11 +10,11 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.yaml.snakeyaml.Yaml; -import com.alibaba.otter.canal.kafka.CanalServerStarter; -import com.alibaba.otter.canal.kafka.producer.KafkaProperties.CanalDestination; -import com.alibaba.otter.canal.kafka.producer.KafkaProperties.Topic; +import com.alibaba.otter.canal.kafka.KafkaProperties.CanalDestination; +import com.alibaba.otter.canal.kafka.KafkaProperties.Topic; import com.alibaba.otter.canal.protocol.ClientIdentity; import com.alibaba.otter.canal.protocol.Message; +import com.alibaba.otter.canal.server.CanalServerStarter; import com.alibaba.otter.canal.server.embedded.CanalServerWithEmbedded; /** @@ -23,20 +23,21 @@ import com.alibaba.otter.canal.server.embedded.CanalServerWithEmbedded; * @author machengyuan 2018-6-11 下午05:30:49 * @version 1.0.0 */ -public class CanalKafkaStarter { +public class CanalKafkaStarter implements CanalServerStarter { - private static final String CLASSPATH_URL_PREFIX = "classpath:"; - private static final Logger logger = LoggerFactory.getLogger(CanalKafkaStarter.class); + private static final Logger logger = LoggerFactory.getLogger(CanalKafkaStarter.class); - private volatile static boolean running = false; + private static final String CLASSPATH_URL_PREFIX = "classpath:"; - private static ExecutorService executorService; + private volatile boolean running = false; - private static CanalKafkaProducer canalKafkaProducer; + private ExecutorService executorService; - private static KafkaProperties kafkaProperties; + private CanalKafkaProducer canalKafkaProducer; - public static void init() { + private KafkaProperties kafkaProperties; + + public void init() { try { logger.info("## load kafka configurations"); String conf = System.getProperty("kafka.conf", "classpath:kafka.yml"); @@ -96,11 +97,9 @@ public class CanalKafkaStarter { } } - private static void worker(CanalDestination destination) { + private void worker(CanalDestination destination) { while (!running) ; - while (!CanalServerStarter.isRunning()) - ; // 等待server启动完成 logger.info("## start the canal consumer: {}.", destination.getCanalDestination()); CanalServerWithEmbedded server = CanalServerWithEmbedded.instance(); ClientIdentity clientIdentity = new ClientIdentity(destination.getCanalDestination(), (short) 1001, ""); @@ -121,7 +120,7 @@ public class CanalKafkaStarter { Message message = server.getWithoutAck(clientIdentity, kafkaProperties.getCanalBatchSize()); // 获取指定数量的数据 long batchId = message.getId(); try { - int size = message.isRaw() ? message.getRawEntries().size() : message.getEntries().size(); + int size = message.isRaw() ? message.getRawEntries().size() : message.getEntries().size(); if (batchId != -1 && size != 0) { if (!StringUtils.isEmpty(destination.getTopic())) { Topic topic = new Topic(); diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java b/server/src/main/java/com/alibaba/otter/canal/kafka/KafkaProperties.java similarity index 98% rename from kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java rename to server/src/main/java/com/alibaba/otter/canal/kafka/KafkaProperties.java index 8e193aac..140ce965 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/KafkaProperties.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.producer; +package com.alibaba.otter.canal.kafka; import java.util.ArrayList; import java.util.HashSet; diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/MessageSerializer.java b/server/src/main/java/com/alibaba/otter/canal/kafka/MessageSerializer.java similarity index 83% rename from kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/MessageSerializer.java rename to server/src/main/java/com/alibaba/otter/canal/kafka/MessageSerializer.java index 02336004..49db5aa1 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/MessageSerializer.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/MessageSerializer.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.producer; +package com.alibaba.otter.canal.kafka; import java.util.List; import java.util.Map; @@ -37,20 +37,20 @@ public class MessageSerializer implements Serializer { List rowEntries = data.getRawEntries(); // message size int messageSize = 0; - messageSize += com.google.protobuf.CodedOutputStream.computeInt64Size(1, data.getId()); + messageSize += CodedOutputStream.computeInt64Size(1, data.getId()); int dataSize = 0; for (int i = 0; i < rowEntries.size(); i++) { - dataSize += com.google.protobuf.CodedOutputStream.computeBytesSizeNoTag(rowEntries.get(i)); + dataSize += CodedOutputStream.computeBytesSizeNoTag(rowEntries.get(i)); } messageSize += dataSize; messageSize += 1 * rowEntries.size(); // packet size int size = 0; - size += com.google.protobuf.CodedOutputStream.computeEnumSize(3, + size += CodedOutputStream.computeEnumSize(3, PacketType.MESSAGES.getNumber()); - size += com.google.protobuf.CodedOutputStream.computeTagSize(5) - + com.google.protobuf.CodedOutputStream.computeRawVarint32Size(messageSize) + size += CodedOutputStream.computeTagSize(5) + + CodedOutputStream.computeRawVarint32Size(messageSize) + messageSize; // build data byte[] body = new byte[size]; @@ -73,7 +73,7 @@ public class MessageSerializer implements Serializer { } CanalPacket.Packet.Builder packetBuilder = CanalPacket.Packet.newBuilder(); - packetBuilder.setType(CanalPacket.PacketType.MESSAGES); + packetBuilder.setType(PacketType.MESSAGES); packetBuilder.setBody(messageBuilder.build().toByteString()); return packetBuilder.build().toByteArray(); } diff --git a/server/src/main/java/com/alibaba/otter/canal/server/CanalServerStarter.java b/server/src/main/java/com/alibaba/otter/canal/server/CanalServerStarter.java new file mode 100644 index 00000000..18f3f7d4 --- /dev/null +++ b/server/src/main/java/com/alibaba/otter/canal/server/CanalServerStarter.java @@ -0,0 +1,12 @@ +package com.alibaba.otter.canal.server; + +/** + * 外部服务如Kafka, RocketMQ启动接口 + * + * @author machengyuan 2018-8-23 下午05:20:29 + * @version 1.0.0 + */ +public interface CanalServerStarter { + + void init(); +} From fc3de4b7fdf5ce74f0f82fbc0ae2be3d855baf73 Mon Sep 17 00:00:00 2001 From: mcy Date: Thu, 23 Aug 2018 17:38:17 +0800 Subject: [PATCH 2/5] =?UTF-8?q?=E6=95=B4=E5=90=88canal.client=E5=92=8Ccana?= =?UTF-8?q?l.kafka-client?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client-launcher/pom.xml | 9 +- .../loader/CanalAdapterKafkaWorker.java | 22 ++-- client/pom.xml | 7 ++ .../client/kafka}/KafkaCanalConnector.java | 25 ++-- .../client/kafka}/KafkaCanalConnectors.java | 2 +- .../client/kafka}/MessageDeserializer.java | 2 +- .../kafka}/running/ClientRunningData.java | 2 +- .../kafka}/running/ClientRunningListener.java | 2 +- .../kafka}/running/ClientRunningMonitor.java | 10 +- .../running/kafka}/AbstractKafkaTest.java | 2 +- .../kafka}/CanalKafkaClientExample.java | 6 +- .../kafka}/KafkaClientRunningTest.java | 6 +- kafka-client/pom.xml | 111 ------------------ kafka-client/src/test/resources/logback.xml | 19 --- pom.xml | 1 - 15 files changed, 54 insertions(+), 172 deletions(-) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/KafkaCanalConnector.java (96%) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/KafkaCanalConnectors.java (97%) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/MessageDeserializer.java (97%) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/running/ClientRunningData.java (92%) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/running/ClientRunningListener.java (86%) rename {kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client => client/src/main/java/com/alibaba/otter/canal/client/kafka}/running/ClientRunningMonitor.java (97%) rename {kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running => client/src/test/java/com/alibaba/otter/canal/client/running/kafka}/AbstractKafkaTest.java (92%) rename {kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running => client/src/test/java/com/alibaba/otter/canal/client/running/kafka}/CanalKafkaClientExample.java (96%) rename {kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running => client/src/test/java/com/alibaba/otter/canal/client/running/kafka}/KafkaClientRunningTest.java (90%) delete mode 100644 kafka-client/pom.xml delete mode 100644 kafka-client/src/test/resources/logback.xml diff --git a/client-launcher/pom.xml b/client-launcher/pom.xml index 952aeb8a..1f63fe73 100644 --- a/client-launcher/pom.xml +++ b/client-launcher/pom.xml @@ -18,17 +18,16 @@ client-adapter.common ${project.version} - com.alibaba.otter canal.client ${project.version} - + - com.alibaba.otter - canal.kafka.client - ${project.version} + org.apache.kafka + kafka-clients + 1.1.1 org.yaml diff --git a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java index 6b35e452..1842cddb 100644 --- a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java +++ b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java @@ -1,19 +1,25 @@ package com.alibaba.otter.canal.client.adapter.loader; -import com.alibaba.otter.canal.client.adapter.CanalOuterAdapter; -import com.alibaba.otter.canal.client.adapter.loader.AbstractCanalAdapterWorker; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnector; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnectors; -import com.alibaba.otter.canal.protocol.Message; -import org.apache.kafka.clients.consumer.CommitFailedException; -import org.apache.kafka.common.errors.WakeupException; - import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.kafka.clients.consumer.CommitFailedException; +import org.apache.kafka.common.errors.WakeupException; + +import com.alibaba.otter.canal.client.adapter.CanalOuterAdapter; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnector; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnectors; +import com.alibaba.otter.canal.protocol.Message; + +/** + * kafka对应的client适配器工作线程 + * + * @author machengyuan 2018-8-19 下午11:30:49 + * @version 1.0.0 + */ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { private KafkaCanalConnector connector; diff --git a/client/pom.xml b/client/pom.xml index 038f2a9f..6318f771 100644 --- a/client/pom.xml +++ b/client/pom.xml @@ -16,6 +16,13 @@ canal.protocol ${project.version} + + + org.apache.kafka + kafka-clients + 1.1.1 + provided + diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java similarity index 96% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java index 3e997e27..794f80ab 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java @@ -1,21 +1,22 @@ -package com.alibaba.otter.canal.kafka.client; +package com.alibaba.otter.canal.client.kafka; + +import java.util.Collections; +import java.util.Properties; +import java.util.concurrent.TimeUnit; -import com.alibaba.otter.canal.common.utils.AddressUtils; -import com.alibaba.otter.canal.common.utils.BooleanMutex; -import com.alibaba.otter.canal.common.zookeeper.ZkClientx; -import com.alibaba.otter.canal.kafka.client.running.ClientRunningData; -import com.alibaba.otter.canal.kafka.client.running.ClientRunningListener; -import com.alibaba.otter.canal.kafka.client.running.ClientRunningMonitor; -import com.alibaba.otter.canal.protocol.Message; -import com.alibaba.otter.canal.protocol.exception.CanalClientException; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.StringDeserializer; -import java.util.Collections; -import java.util.Properties; -import java.util.concurrent.TimeUnit; +import com.alibaba.otter.canal.client.kafka.running.ClientRunningData; +import com.alibaba.otter.canal.client.kafka.running.ClientRunningListener; +import com.alibaba.otter.canal.client.kafka.running.ClientRunningMonitor; +import com.alibaba.otter.canal.common.utils.AddressUtils; +import com.alibaba.otter.canal.common.utils.BooleanMutex; +import com.alibaba.otter.canal.common.zookeeper.ZkClientx; +import com.alibaba.otter.canal.protocol.Message; +import com.alibaba.otter.canal.protocol.exception.CanalClientException; /** * canal kafka 数据操作客户端 diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnectors.java similarity index 97% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnectors.java index 1a2942ab..0634a747 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnectors.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client; +package com.alibaba.otter.canal.client.kafka; /** * canal kafka connectors创建工具类 diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/MessageDeserializer.java similarity index 97% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/MessageDeserializer.java index 70f2fc06..e716d2d6 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/MessageDeserializer.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client; +package com.alibaba.otter.canal.client.kafka; import java.util.Map; diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningData.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningData.java similarity index 92% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningData.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningData.java index 7bbbabd7..e528f9ee 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningData.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningData.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.kafka.running; /** * client running状态信息 diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningListener.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningListener.java similarity index 86% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningListener.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningListener.java index d49cf79f..016f7f10 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningListener.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningListener.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.kafka.running; /** * client running状态信息 diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningMonitor.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningMonitor.java similarity index 97% rename from kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningMonitor.java rename to client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningMonitor.java index 0f61211e..9a84d81e 100644 --- a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningMonitor.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/running/ClientRunningMonitor.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.kafka.running; import java.text.MessageFormat; import java.util.Random; @@ -55,15 +55,15 @@ public class ClientRunningMonitor extends AbstractCanalLifeCycle { } private static final Logger logger = LoggerFactory.getLogger(ClientRunningMonitor.class); - private ZkClientx zkClient; + private ZkClientx zkClient; private String topic; - private ClientRunningData clientData; + private ClientRunningData clientData; private IZkDataListener dataListener; - private BooleanMutex mutex = new BooleanMutex(false); + private BooleanMutex mutex = new BooleanMutex(false); private volatile boolean release = false; private volatile ClientRunningData activeData; private ScheduledExecutorService delayExector = Executors.newScheduledThreadPool(1); - private ClientRunningListener listener; + private ClientRunningListener listener; private int delayTime = 5; private static Integer virtualPort; diff --git a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java similarity index 92% rename from kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java rename to client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java index 22dda072..473baf67 100644 --- a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java +++ b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.running.kafka; import org.junit.Assert; diff --git a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/CanalKafkaClientExample.java b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java similarity index 96% rename from kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/CanalKafkaClientExample.java rename to client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java index f4a0afc5..0a53a1df 100644 --- a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/CanalKafkaClientExample.java +++ b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.running.kafka; import java.util.concurrent.TimeUnit; @@ -7,8 +7,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.util.Assert; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnector; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnectors; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnector; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnectors; import com.alibaba.otter.canal.protocol.Message; /** diff --git a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/KafkaClientRunningTest.java b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java similarity index 90% rename from kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/KafkaClientRunningTest.java rename to client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java index d8ff6d30..ebc4d549 100644 --- a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/KafkaClientRunningTest.java +++ b/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java @@ -1,4 +1,4 @@ -package com.alibaba.otter.canal.kafka.client.running; +package com.alibaba.otter.canal.client.running.kafka; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -9,8 +9,8 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnector; -import com.alibaba.otter.canal.kafka.client.KafkaCanalConnectors; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnector; +import com.alibaba.otter.canal.client.kafka.KafkaCanalConnectors; import com.alibaba.otter.canal.protocol.Message; /** diff --git a/kafka-client/pom.xml b/kafka-client/pom.xml deleted file mode 100644 index 42424ddd..00000000 --- a/kafka-client/pom.xml +++ /dev/null @@ -1,111 +0,0 @@ - - - 4.0.0 - - canal - com.alibaba.otter - 1.1.0-SNAPSHOT - ../pom.xml - - com.alibaba.otter - canal.kafka.client - jar - canal kafka client module for otter ${project.version} - - - - - com.alibaba.otter - canal.protocol - ${project.version} - - - org.apache.kafka - kafka-clients - 1.1.1 - - - - - junit - junit - - - - - - dev - - true - - env - !javadoc - - - - - - javadoc - - - env - javadoc - - - - - - org.apache.maven.plugins - maven-javadoc-plugin - 2.9.1 - - - attach-javadocs - package - - jar - - - - - true - public - true -
${project.artifactId}-${project.version}
-
${project.artifactId}-${project.version}
- ${project.artifactId}-${project.version} - - https://github.com/alibaba/canal - - ${project.build.directory}/apidocs/apidocs/${project.version} -
-
- - org.apache.maven.plugins - maven-scm-publish-plugin - 1.0-beta-2 - - - attach-javadocs - package - - publish-scm - - - - - ${project.build.directory}/scmpublish - Publishing javadoc for ${project.artifactId}:${project.version} - ${project.build.directory}/apidocs - true - scm:git:git@github.com:alibaba/canal.git - gh-pages - - -
-
-
-
-
diff --git a/kafka-client/src/test/resources/logback.xml b/kafka-client/src/test/resources/logback.xml deleted file mode 100644 index 81fa0712..00000000 --- a/kafka-client/src/test/resources/logback.xml +++ /dev/null @@ -1,19 +0,0 @@ - - - - - - %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{56} - %msg%n - - - - - - - - - - - - - \ No newline at end of file diff --git a/pom.xml b/pom.xml index 64588a28..71ff9fd5 100644 --- a/pom.xml +++ b/pom.xml @@ -117,7 +117,6 @@ client deployer example - kafka-client prometheus client-adapter client-launcher From 9e3d4a7ffe5371557d79de1899fb68152a1a1d83 Mon Sep 17 00:00:00 2001 From: mcy Date: Thu, 23 Aug 2018 17:43:53 +0800 Subject: [PATCH 3/5] =?UTF-8?q?HBase=E8=AF=BB=E9=85=8D=E7=BD=AE=E7=A9=BA?= =?UTF-8?q?=E8=A1=8C=E7=9A=84bug?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../canal/client/adapter/hbase/config/MappingConfigLoader.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 f00e6282..929e9d6a 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 @@ -45,7 +45,7 @@ public class MappingConfigLoader { continue; } c = c.trim(); - if (c.startsWith("#")) { + if (c.equals("") || c.startsWith("#")) { continue; } From b7de1e50705fc196d2261a8d983f0228bd86d960 Mon Sep 17 00:00:00 2001 From: mcy Date: Fri, 24 Aug 2018 09:14:52 +0800 Subject: [PATCH 4/5] =?UTF-8?q?=E5=88=A0=E9=99=A4kafka=20kafka-client?= =?UTF-8?q?=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kafka-client/pom.xml | 0 kafka/pom.xml | 0 2 files changed, 0 insertions(+), 0 deletions(-) delete mode 100644 kafka-client/pom.xml delete mode 100644 kafka/pom.xml diff --git a/kafka-client/pom.xml b/kafka-client/pom.xml deleted file mode 100644 index e69de29b..00000000 diff --git a/kafka/pom.xml b/kafka/pom.xml deleted file mode 100644 index e69de29b..00000000 From 6b94f396ffaca3dfcaef5fb6af4c152982054399 Mon Sep 17 00:00:00 2001 From: mcy Date: Fri, 24 Aug 2018 10:52:04 +0800 Subject: [PATCH 5/5] =?UTF-8?q?=E7=A7=BB=E9=99=A4kafka=20client=E5=9F=BA?= =?UTF-8?q?=E4=BA=8Ezk=E7=9A=84HA=E6=9C=BA=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/support/CanalClientConfig.java | 4 +- .../loader/CanalAdapterKafkaWorker.java | 5 +- .../adapter/loader/CanalAdapterLoader.java | 89 ++++++++++--------- .../src/main/resources/canal-client.yml | 19 ++-- .../client/kafka/KafkaCanalConnector.java | 78 ++++++++-------- 5 files changed, 99 insertions(+), 96 deletions(-) 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 4f8ff396..f6fa7707 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 @@ -20,9 +20,9 @@ public class CanalClientConfig { private String bootstrapServers; - private List kafkaTopics = new ArrayList<>(); + private List kafkaTopics; - private List canalInstances = new ArrayList<>(); + private List canalInstances; public String getCanalServerHost() { return canalServerHost; diff --git a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java index 1842cddb..5bccb51a 100644 --- a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java +++ b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterKafkaWorker.java @@ -91,7 +91,7 @@ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { private void process() { while (!running) ; - ExecutorService executor = Executors.newFixedThreadPool(1); + ExecutorService executor = Executors.newSingleThreadExecutor(); final AtomicBoolean executing = new AtomicBoolean(true); while (running) { try { @@ -142,7 +142,8 @@ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { } }); - while (executing.get()) { // keeping kafka client active + // 间隔一段时间ack一次, 防止因超时未响应切换到另外台客户端 + while (executing.get()) { connector.ack(); Thread.sleep(500); } diff --git a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterLoader.java b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterLoader.java index b0fe35c5..98973f11 100644 --- a/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterLoader.java +++ b/client-launcher/src/main/java/com/alibaba/otter/canal/client/adapter/loader/CanalAdapterLoader.java @@ -26,13 +26,13 @@ public class CanalAdapterLoader { private static final Logger logger = LoggerFactory.getLogger(CanalAdapterLoader.class); - private CanalClientConfig canalClientConfig; + private CanalClientConfig canalClientConfig; private Map canalWorkers = new HashMap<>(); private Map canalKafkaWorkers = new HashMap<>(); - private ExtensionLoader loader; + private ExtensionLoader loader; public CanalAdapterLoader(CanalClientConfig canalClientConfig){ this.canalClientConfig = canalClientConfig; @@ -43,7 +43,7 @@ public class CanalAdapterLoader { */ public void init() { // canal instances 和 kafka topics 配置不能同时为空 - if (canalClientConfig.getCanalInstances().isEmpty() && canalClientConfig.getKafkaTopics().isEmpty()) { + if (canalClientConfig.getCanalInstances() == null && canalClientConfig.getKafkaTopics() == null) { throw new RuntimeException("Blank config property: canalInstances or canalKafkaTopics"); } @@ -58,56 +58,61 @@ public class CanalAdapterLoader { } String zkHosts = this.canalClientConfig.getZookeeperHosts(); - if (zkHosts == null && sa == null) { - throw new RuntimeException("Blank config property: canalServerHost or zookeeperHosts"); - } + // if (zkHosts == null && sa == null) { + // throw new RuntimeException("Blank config property: canalServerHost or + // zookeeperHosts"); + // } // 初始化canal-client的适配器 - for (CanalClientConfig.CanalInstance instance : canalClientConfig.getCanalInstances()) { - List> canalOuterAdapterGroups = new ArrayList<>(); + if (canalClientConfig.getCanalInstances() != null) { + for (CanalClientConfig.CanalInstance instance : canalClientConfig.getCanalInstances()) { + List> canalOuterAdapterGroups = new ArrayList<>(); - for (CanalClientConfig.AdapterGroup connectorGroup : instance.getAdapterGroups()) { - List canalOutConnectors = new ArrayList<>(); - for (CanalOuterAdapterConfiguration c : connectorGroup.getOutAdapters()) { - loadConnector(c, canalOutConnectors); + for (CanalClientConfig.AdapterGroup connectorGroup : instance.getAdapterGroups()) { + List canalOutConnectors = new ArrayList<>(); + for (CanalOuterAdapterConfiguration c : connectorGroup.getOutAdapters()) { + loadConnector(c, canalOutConnectors); + } + canalOuterAdapterGroups.add(canalOutConnectors); } - canalOuterAdapterGroups.add(canalOutConnectors); + CanalAdapterWorker worker; + if (zkHosts != null) { + worker = new CanalAdapterWorker(instance.getInstance(), zkHosts, canalOuterAdapterGroups); + } else { + worker = new CanalAdapterWorker(instance.getInstance(), sa, canalOuterAdapterGroups); + } + canalWorkers.put(instance.getInstance(), worker); + worker.start(); + logger.info("Start adapter for canal instance: {} succeed", instance.getInstance()); } - CanalAdapterWorker worker; - if (zkHosts != null) { - worker = new CanalAdapterWorker(instance.getInstance(), zkHosts, canalOuterAdapterGroups); - } else { - worker = new CanalAdapterWorker(instance.getInstance(), sa, canalOuterAdapterGroups); - } - canalWorkers.put(instance.getInstance(), worker); - worker.start(); - logger.info("Start adapter for canal instance: {} succeed", instance.getInstance()); } // 初始化canal-client-kafka的适配器 - for (CanalClientConfig.KafkaTopic kafkaTopic : canalClientConfig.getKafkaTopics()) { - for (CanalClientConfig.Group group : kafkaTopic.getGroups()) { - List> canalOuterAdapterGroups = new ArrayList<>(); + if (canalClientConfig.getKafkaTopics() != null) { + for (CanalClientConfig.KafkaTopic kafkaTopic : canalClientConfig.getKafkaTopics()) { + for (CanalClientConfig.Group group : kafkaTopic.getGroups()) { + List> canalOuterAdapterGroups = new ArrayList<>(); - List canalOuterAdapters = new ArrayList<>(); + List canalOuterAdapters = new ArrayList<>(); - for (CanalOuterAdapterConfiguration config : group.getOutAdapters()) { - // for (CanalOuterAdapterConfiguration config : adaptor.getOutAdapters()) { - loadConnector(config, canalOuterAdapters); - // } + for (CanalOuterAdapterConfiguration config : group.getOutAdapters()) { + // for (CanalOuterAdapterConfiguration config : adaptor.getOutAdapters()) { + loadConnector(config, canalOuterAdapters); + // } + } + canalOuterAdapterGroups.add(canalOuterAdapters); + + // String zkServers = canalClientConfig.getZookeeperHosts(); + CanalAdapterKafkaWorker canalKafkaWorker = new CanalAdapterKafkaWorker(zkHosts, + canalClientConfig.getBootstrapServers(), + kafkaTopic.getTopic(), + group.getGroupId(), + canalOuterAdapterGroups); + canalKafkaWorkers.put(kafkaTopic.getTopic() + "-" + group.getGroupId(), canalKafkaWorker); + canalKafkaWorker.start(); + logger.info("Start adapter for canal-client kafka topic: {} succeed", + kafkaTopic.getTopic() + "-" + group.getGroupId()); } - canalOuterAdapterGroups.add(canalOuterAdapters); - - String zkServers = canalClientConfig.getZookeeperHosts(); - CanalAdapterKafkaWorker canalKafkaWorker = new CanalAdapterKafkaWorker(zkServers, - canalClientConfig.getBootstrapServers(), - kafkaTopic.getTopic(), - group.getGroupId(), - canalOuterAdapterGroups); - canalKafkaWorkers.put(kafkaTopic.getTopic() + "-" + group.getGroupId(), canalKafkaWorker); - canalKafkaWorker.start(); - logger.info("Start adapter for canal-client kafka topic: {} succeed", - kafkaTopic.getTopic() + "-" + group.getGroupId()); } } } diff --git a/client-launcher/src/main/resources/canal-client.yml b/client-launcher/src/main/resources/canal-client.yml index 2412e38d..12e016a7 100644 --- a/client-launcher/src/main/resources/canal-client.yml +++ b/client-launcher/src/main/resources/canal-client.yml @@ -1,6 +1,6 @@ canalServerHost: 127.0.0.1:11111 -#zookeeperHosts: 127.0.0.1:2181 -#bootstrapServers: kafka1.mytest.com:9092,kafka2.mytest.com:9092 +#zookeeperHosts: slave1:2181 +#bootstrapServers: slave1:6667,slave2:6667 canalInstances: - instance: example @@ -10,13 +10,12 @@ canalInstances: - name: hbase hosts: slave1:2181 properties: {znodeParent: "/hbase-unsecure"} - #kafkaTopics: -#- topic: devmysql4308 +#- topic: example # groups: -# - groupId: devmysql4308_es -# adapters: -# - name: es -# hosts: -# zkHosts: -# properties: {clusterName: es-service-test} +# - groupId: example_g1 +# outAdapters: +# - name: logger +# - name: hbase +# hosts: slave1:2181 +# properties: {znodeParent: "/hbase-unsecure"} diff --git a/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java index 794f80ab..c5b2f8fc 100644 --- a/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaCanalConnector.java @@ -10,10 +10,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.StringDeserializer; import com.alibaba.otter.canal.client.kafka.running.ClientRunningData; -import com.alibaba.otter.canal.client.kafka.running.ClientRunningListener; -import com.alibaba.otter.canal.client.kafka.running.ClientRunningMonitor; import com.alibaba.otter.canal.common.utils.AddressUtils; -import com.alibaba.otter.canal.common.utils.BooleanMutex; import com.alibaba.otter.canal.common.zookeeper.ZkClientx; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.exception.CanalClientException; @@ -27,16 +24,16 @@ import com.alibaba.otter.canal.protocol.exception.CanalClientException; public class KafkaCanalConnector { private KafkaConsumer kafkaConsumer; - private String topic; - private Integer partition; - private Properties properties; - private ClientRunningMonitor runningMonitor; // 运行控制 - private ZkClientx zkClientx; - private BooleanMutex mutex = new BooleanMutex(false); - private volatile boolean connected = false; - private volatile boolean running = false; + private String topic; + private Integer partition; + private Properties properties; + // private ClientRunningMonitor runningMonitor; // 运行控制 + // private BooleanMutex mutex = new BooleanMutex(false); + private ZkClientx zkClientx; + private volatile boolean connected = false; + private volatile boolean running = false; - public KafkaCanalConnector(String zkServers, String servers, String topic, Integer partition, String groupId) { + public KafkaCanalConnector(String zkServers, String servers, String topic, Integer partition, String groupId){ this.topic = topic; this.partition = partition; @@ -45,7 +42,7 @@ public class KafkaCanalConnector { properties.put("group.id", groupId); properties.put("enable.auto.commit", false); properties.put("auto.commit.interval.ms", "1000"); - properties.put("auto.offset.reset", "latest"); //如果没有offset则从最后的offset开始读 + properties.put("auto.offset.reset", "latest"); // 如果没有offset则从最后的offset开始读 properties.put("request.timeout.ms", "40000"); // 必须大于session.timeout.ms的设置 properties.put("session.timeout.ms", "30000"); // 默认为30秒 properties.put("max.poll.records", "1"); // 所以一次只取一条数据 @@ -59,19 +56,19 @@ public class KafkaCanalConnector { clientData.setGroupId(groupId); clientData.setAddress(AddressUtils.getHostIp()); - runningMonitor = new ClientRunningMonitor(); - runningMonitor.setTopic(topic); - runningMonitor.setZkClient(zkClientx); - runningMonitor.setClientData(clientData); - runningMonitor.setListener(new ClientRunningListener() { - public void processActiveEnter() { - mutex.set(true); - } - - public void processActiveExit() { - mutex.set(false); - } - }); + // runningMonitor = new ClientRunningMonitor(); + // runningMonitor.setTopic(topic); + // runningMonitor.setZkClient(zkClientx); + // runningMonitor.setClientData(clientData); + // runningMonitor.setListener(new ClientRunningListener() { + // public void processActiveEnter() { + // mutex.set(true); + // } + // + // public void processActiveExit() { + // mutex.set(false); + // } + // }); } } @@ -96,11 +93,11 @@ public class KafkaCanalConnector { return; } - if (runningMonitor != null) { - if (!runningMonitor.isStart()) { - runningMonitor.start(); - } - } + // if (runningMonitor != null) { + // if (!runningMonitor.isStart()) { + // runningMonitor.start(); + // } + // } connected = true; @@ -116,9 +113,9 @@ public class KafkaCanalConnector { kafkaConsumer.close(); connected = false; - if (runningMonitor.isStart()) { - runningMonitor.stop(); - } + // if (runningMonitor.isStart()) { + // runningMonitor.stop(); + // } } private void waitClientRunning() { @@ -129,12 +126,12 @@ public class KafkaCanalConnector { } running = true; - mutex.get();// 阻塞等待 + // mutex.get();// 阻塞等待 } else { // 单机模式直接设置为running running = true; } - } catch (InterruptedException e) { + } catch (Exception e) { Thread.currentThread().interrupt(); throw new CanalClientException(e); } @@ -142,7 +139,8 @@ public class KafkaCanalConnector { public boolean checkValid() { if (zkClientx != null) { - return mutex.state(); + // return mutex.state(); + return true; } else { return true;// 默认都放过 } @@ -235,9 +233,9 @@ public class KafkaCanalConnector { public void stopRunning() { if (running) { running = false; // 设置为非running状态 - if (!mutex.state()) { - mutex.set(true); // 中断阻塞 - } + // if (!mutex.state()) { + // mutex.set(true); // 中断阻塞 + // } } } }