Merge pull request #1057 from rewerma/master
kafka canal-adapter 分布式开关的bug
This commit is contained in:
@@ -1 +0,0 @@
|
||||
/bin/
|
||||
+7
-2
@@ -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) {
|
||||
|
||||
+7
@@ -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);
|
||||
|
||||
@@ -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"}
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user