diff --git a/client-adapter/example/.gitignore b/client-adapter/example/.gitignore deleted file mode 100644 index ae3c1726..00000000 --- a/client-adapter/example/.gitignore +++ /dev/null @@ -1 +0,0 @@ -/bin/ diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java index e74641ca..41664c7c 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterKafkaWorker.java @@ -35,7 +35,7 @@ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { } @Override - protected void closeConnection(){ + protected void closeConnection() { connector.stopRunning(); } @@ -45,6 +45,7 @@ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { ; while (running) { try { + syncSwitch.get(canalDestination); logger.info("=============> Start to connect topic: {} <=============", this.topic); connector.connect(); logger.info("=============> Start to subscribe topic: {} <=============", this.topic); @@ -52,7 +53,11 @@ public class CanalAdapterKafkaWorker extends AbstractCanalAdapterWorker { logger.info("=============> Subscribe topic: {} succeed <=============", this.topic); while (running) { try { - syncSwitch.get(canalDestination); + Boolean status = syncSwitch.status(canalDestination); + if (status != null && !status) { + connector.disconnect(); + break; + } List messages; if (!flatMessage) { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java index db80d602..5f5b058c 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java @@ -43,6 +43,7 @@ public class CanalAdapterRocketMQWorker extends AbstractCanalAdapterWorker { ; while (running) { try { + syncSwitch.get(canalDestination); logger.info("=============> Start to connect topic: {} <=============", this.topic); connector.connect(); logger.info("=============> Start to subscribe topic: {}<=============", this.topic); @@ -50,6 +51,12 @@ public class CanalAdapterRocketMQWorker extends AbstractCanalAdapterWorker { logger.info("=============> Subscribe topic: {} succeed<=============", this.topic); while (running) { try { + Boolean status = syncSwitch.status(canalDestination); + if (status != null && !status) { + connector.disconnect(); + break; + } + List messages; if (!flatMessage) { messages = connector.getListWithoutAck(100L, TimeUnit.MILLISECONDS); diff --git a/client-launcher/src/main/resources/canal-client.yml b/client-launcher/src/main/resources/canal-client.yml deleted file mode 100644 index 211296b5..00000000 --- a/client-launcher/src/main/resources/canal-client.yml +++ /dev/null @@ -1,24 +0,0 @@ -#canalServerHost: 127.0.0.1:11111 -#zookeeperHosts: slave1:2181 -bootstrapServers: slave1:6667,slave2:6667 #or rocketmq nameservers:host1:9876;host2:9876 -flatMessage: false - -#canalInstances: -#- instance: example -# adapterGroups: -# - outAdapters: -# - name: logger -# - name: hbase -# hosts: slave1:2181 -# properties: {znodeParent: "/hbase-unsecure"} - -mqTopics: -- mqMode: rocketmq - topic: example - groups: - - groupId: example - 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 6439b1b0..e9282498 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 @@ -87,9 +87,11 @@ public class KafkaCanalConnector implements CanalMQConnector { public void disconnect() { if (kafkaConsumer != null) { kafkaConsumer.close(); + kafkaConsumer = null; } if (kafkaConsumer2 != null) { kafkaConsumer2.close(); + kafkaConsumer2 = null; } connected = false;