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