From 330da50a5ecdfd1bb0d9fb1b064e2918d950e392 Mon Sep 17 00:00:00 2001 From: rewerma Date: Tue, 12 Jun 2018 10:23:50 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9kafka=E9=85=8D=E7=BD=AE?= =?UTF-8?q?=E9=A1=B9=E5=B1=9E=E6=80=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../otter/canal/kafka/producer/CanalKafkaStarter.java | 6 +++--- .../otter/canal/kafka/producer/KafkaProperties.java | 10 +++++----- kafka/src/main/resources/kafka.yml | 2 +- 3 files changed, 9 insertions(+), 9 deletions(-) diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java b/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java index 7514c5ce..3ef05489 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java +++ b/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/CanalKafkaStarter.java @@ -91,9 +91,9 @@ public class CanalKafkaStarter { private static void worker(Topic topic) { while (!running) ; while (!CanalServerStarter.isRunning()) ; //等待server启动完成 - logger.info("## start the canal consumer: {}.", topic.getDestination()); + logger.info("## start the canal consumer: {}.", topic.getCanalDestination()); CanalServerWithEmbedded server = CanalServerWithEmbedded.instance(); - ClientIdentity clientIdentity = new ClientIdentity(topic.getDestination(), (short) 1001, ""); + ClientIdentity clientIdentity = new ClientIdentity(topic.getCanalDestination(), (short) 1001, ""); while (running) { try { if (!server.getCanalInstances().containsKey(clientIdentity.getDestination())) { @@ -105,7 +105,7 @@ public class CanalKafkaStarter { continue; } server.subscribe(clientIdentity); - logger.info("## the canal consumer {} is running now ......", topic.getDestination()); + logger.info("## the canal consumer {} is running now ......", topic.getCanalDestination()); while (running) { Message message = server.getWithoutAck(clientIdentity, 5 * 1024); // 获取指定数量的数据 diff --git a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java b/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java index f70ac2c4..f7362a78 100644 --- a/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java +++ b/kafka/src/main/java/com/alibaba/otter/canal/kafka/producer/KafkaProperties.java @@ -21,7 +21,7 @@ public class KafkaProperties { public static class Topic { private String topic; private Integer partition; - private String destination; + private String canalDestination; public String getTopic() { return topic; @@ -39,12 +39,12 @@ public class KafkaProperties { this.partition = partition; } - public String getDestination() { - return destination; + public String getCanalDestination() { + return canalDestination; } - public void setDestination(String destination) { - this.destination = destination; + public void setCanalDestination(String canalDestination) { + this.canalDestination = canalDestination; } } diff --git a/kafka/src/main/resources/kafka.yml b/kafka/src/main/resources/kafka.yml index 3a846c61..7ca1bb7f 100644 --- a/kafka/src/main/resources/kafka.yml +++ b/kafka/src/main/resources/kafka.yml @@ -7,6 +7,6 @@ bufferMemory: 33554432 topics: - topic: example partition: - destination: example + canalDestination: example