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); // 中断阻塞 + // } } } }