From e7846085fc60f382ce6f6bf899e646c3d03ffa8f Mon Sep 17 00:00:00 2001 From: rewerma Date: Tue, 12 Jun 2018 19:25:03 +0800 Subject: [PATCH] =?UTF-8?q?kafka=20=E5=AE=A2=E6=88=B7=E7=AB=AF=E6=B6=88?= =?UTF-8?q?=E8=B4=B9=E5=AE=8C=E6=88=90=E5=9F=BA=E6=9C=AC=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kafka-client/pom.xml | 110 +++++++++++ .../kafka/client/KafkaCanalConnector.java | 172 ++++++++++++++++++ .../kafka/client/KafkaCanalConnectors.java | 30 +++ .../kafka/client/MessageDeserializer.java | 55 ++++++ .../client/running/AbstractKafkaTest.java | 19 ++ .../client/running/ClientRunningTest.java | 42 +++++ kafka-client/src/test/resources/logback.xml | 19 ++ kafka/pom.xml | 6 - kafka/src/main/resources/kafka.yml | 5 +- pom.xml | 1 + 10 files changed, 452 insertions(+), 7 deletions(-) create mode 100644 kafka-client/pom.xml create mode 100644 kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java create mode 100644 kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java create mode 100644 kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java create mode 100644 kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java create mode 100644 kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningTest.java create mode 100644 kafka-client/src/test/resources/logback.xml diff --git a/kafka-client/pom.xml b/kafka-client/pom.xml new file mode 100644 index 00000000..0efd13d4 --- /dev/null +++ b/kafka-client/pom.xml @@ -0,0 +1,110 @@ + + + 4.0.0 + + canal + com.alibaba.otter + 1.0.26-SNAPSHOT + + com.alibaba.otter + canal.kafka.client + jar + canal kafka client module for otter ${project.version} + + + + + com.alibaba.otter + canal.client + ${project.version} + + + org.apache.kafka + kafka-clients + 0.10.0.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 + + +
+
+
+
+
\ No newline at end of file diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java new file mode 100644 index 00000000..fc749fec --- /dev/null +++ b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnector.java @@ -0,0 +1,172 @@ +package com.alibaba.otter.canal.kafka.client; + +import com.alibaba.otter.canal.client.CanalConnector; +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.errors.WakeupException; +import org.apache.kafka.common.serialization.StringDeserializer; + +import java.util.Collections; +import java.util.ConcurrentModificationException; +import java.util.Properties; +import java.util.concurrent.TimeUnit; + +/** + * canal kafka 数据操作客户端 + * + * @author machengyuan @ 2018-6-12 + * @version 1.0.0 + */ +public class KafkaCanalConnector implements CanalConnector { + + private KafkaConsumer kafkaConsumer; + + private String topic; + + private Integer partition; + + + private Properties properties; + + public KafkaCanalConnector(String servers, String topic, Integer partition, String groupId) { + this.topic = topic; + this.partition = partition; + + properties = new Properties(); + properties.put("bootstrap.servers", servers); + properties.put("group.id", groupId); + properties.put("enable.auto.commit", false); + properties.put("auto.commit.interval.ms", 1000); + properties.put("auto.offset.reset", "latest"); //earliest //如果没有offset则从最后的offset开始读 + properties.put("request.timeout.ms", 600000); + properties.put("offsets.commit.timeout.ms", 300000); + properties.put("session.timeout.ms", 30000); + properties.put("max.poll.records", 1); //一次只取一条message + properties.put("key.deserializer", StringDeserializer.class.getName()); + properties.put("value.deserializer", MessageDeserializer.class.getName()); + } + + @Override + public void connect() throws CanalClientException { + kafkaConsumer = new KafkaConsumer(properties); + } + + @Override + public void disconnect() throws CanalClientException { + if (kafkaConsumer != null) { + try { + kafkaConsumer.close(); + } catch (ConcurrentModificationException e) { + kafkaConsumer.wakeup(); //通过wakeup异常间接关闭consumer + } + } + } + + @Override + public boolean checkValid() throws CanalClientException { + return true; + } + + @Override + public void subscribe(String filter) throws CanalClientException { + try { + if (kafkaConsumer == null) { + throw new CanalClientException("connect the kafka first before subscribe"); + } + if (partition == null) { + kafkaConsumer.subscribe(Collections.singletonList(topic)); + } else { + kafkaConsumer.subscribe(Collections.singletonList(topic)); + TopicPartition topicPartition = new TopicPartition(topic, partition); + kafkaConsumer.assign(Collections.singletonList(topicPartition)); + } + } catch (WakeupException e) { + closeByWakeupException(e); + } catch (Exception e) { + throw new CanalClientException(e); + } + } + + @Override + public void subscribe() throws CanalClientException { + subscribe(null); + } + + @Override + public void unsubscribe() throws CanalClientException { + try { + kafkaConsumer.unsubscribe(); + } catch (WakeupException e) { + closeByWakeupException(e); + } catch (Exception e) { + throw new CanalClientException(e); + } + } + + @Override + public Message get(int batchSize) throws CanalClientException { + return get(batchSize, 100L, TimeUnit.MILLISECONDS); + } + + @Override + public Message get(int batchSize, Long timeout, TimeUnit unit) throws CanalClientException { + Message message = getWithoutAck(batchSize, timeout, unit); + this.ack(1); + return message; + } + + @Override + public Message getWithoutAck(int batchSize) throws CanalClientException { + return getWithoutAck(batchSize, 100L, TimeUnit.MILLISECONDS); + } + + @Override + public Message getWithoutAck(int batchSize, Long timeout, TimeUnit unit) throws CanalClientException { + try { + ConsumerRecords records = kafkaConsumer.poll(unit.toMillis(timeout)); //基于配置,一次最多只能poll到一条Msg + + if (!records.isEmpty()) { + return records.iterator().next().value(); + } + } catch (WakeupException e) { + closeByWakeupException(e); + } catch (Exception e) { + throw new CanalClientException(e); + } + return null; + } + + @Override + public void ack(long batchId) throws CanalClientException { + try { + kafkaConsumer.commitSync(); + } catch (WakeupException e) { + closeByWakeupException(e); + } catch (Exception e) { + throw new CanalClientException(e); + } + } + + @Override + public void rollback(long batchId) throws CanalClientException { + + } + + @Override + public void rollback() throws CanalClientException { + + } + + @Override + public void stopRunning() throws CanalClientException { + + } + + private void closeByWakeupException(WakeupException e) { + kafkaConsumer.close(); + throw e; + } +} diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java new file mode 100644 index 00000000..f0d4dea6 --- /dev/null +++ b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/KafkaCanalConnectors.java @@ -0,0 +1,30 @@ +package com.alibaba.otter.canal.kafka.client; + +import com.alibaba.otter.canal.client.CanalConnector; + +public class KafkaCanalConnectors { + /** + * 创建kafka客户端链接 + * + * @param servers + * @param topic + * @param partition + * @param groupId + * @return + */ + public static CanalConnector newKafkaConnector(String servers, String topic, Integer partition, String groupId) { + return new KafkaCanalConnector(servers, topic, partition, groupId); + } + + /** + * 创建kafka客户端链接 + * + * @param servers + * @param topic + * @param groupId + * @return + */ + public static CanalConnector newKafkaConnector(String servers, String topic, String groupId) { + return new KafkaCanalConnector(servers, topic, null, groupId); + } +} diff --git a/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java new file mode 100644 index 00000000..c321d076 --- /dev/null +++ b/kafka-client/src/main/java/com/alibaba/otter/canal/kafka/client/MessageDeserializer.java @@ -0,0 +1,55 @@ +package com.alibaba.otter.canal.kafka.client; + +import com.alibaba.otter.canal.protocol.CanalEntry; +import com.alibaba.otter.canal.protocol.CanalPacket; +import com.alibaba.otter.canal.protocol.Message; +import com.alibaba.otter.canal.protocol.exception.CanalClientException; +import com.google.protobuf.ByteString; +import org.apache.kafka.common.serialization.Deserializer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Map; + +public class MessageDeserializer implements Deserializer { + private static Logger logger = LoggerFactory.getLogger(MessageDeserializer.class); + + @Override + public void configure(Map configs, boolean isKey) { + } + + @Override + public Message deserialize(String topic, byte[] data) { + try { + if (data == null) + return null; + else { + CanalPacket.Packet p = CanalPacket.Packet.parseFrom(data); + switch (p.getType()) { + case MESSAGES: { + if (!p.getCompression().equals(CanalPacket.Compression.NONE)) { + throw new CanalClientException("compression is not supported in this connector"); + } + + CanalPacket.Messages messages = CanalPacket.Messages.parseFrom(p.getBody()); + Message result = new Message(messages.getBatchId()); + for (ByteString byteString : messages.getMessagesList()) { + result.addEntry(CanalEntry.Entry.parseFrom(byteString)); + } + return result; + } + default: + break; + } + } + } catch (Exception e) { + logger.error("Error when deserializing byte[] to message "); + } + return null; + } + + @Override + public void close() { + // nothing to do + } +} \ No newline at end of file diff --git a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java b/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java new file mode 100644 index 00000000..2a9ce466 --- /dev/null +++ b/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/AbstractKafkaTest.java @@ -0,0 +1,19 @@ +package com.alibaba.otter.canal.kafka.client.running; + +import org.junit.Assert; + +public class AbstractKafkaTest { + + protected String topic = "example"; + protected Integer partition = null; + protected String groupId = "g1"; + protected String servers = "slave1.test.apitops.com:6667,slave2.test.apitops.com:6667,slave3.test.apitops.com:6667"; + + public void sleep(long time) { + try { + Thread.sleep(time); + } catch (InterruptedException e) { + Assert.fail(e.getMessage()); + } + } +} diff --git a/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningTest.java b/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningTest.java new file mode 100644 index 00000000..c968f0f3 --- /dev/null +++ b/kafka-client/src/test/java/com/alibaba/otter/canal/kafka/client/running/ClientRunningTest.java @@ -0,0 +1,42 @@ +package com.alibaba.otter.canal.kafka.client.running; + +import com.alibaba.otter.canal.client.CanalConnector; +import com.alibaba.otter.canal.kafka.client.KafkaCanalConnectors; +import com.alibaba.otter.canal.protocol.Message; +import org.junit.Test; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +public class ClientRunningTest extends AbstractKafkaTest { + + private boolean running = true; + + @Test + public void testKafkaConsumer() { + final ExecutorService executor = Executors.newFixedThreadPool(1); + + final CanalConnector kafkaCanalConnector = KafkaCanalConnectors.newKafkaConnector(servers, topic, partition, groupId); + kafkaCanalConnector.connect(); + kafkaCanalConnector.subscribe(); + + executor.submit(new Runnable() { + @Override + public void run() { + while (running) { + Message message = kafkaCanalConnector.getWithoutAck(1); + if (message != null) { + System.out.println(message); + sleep(40000); + } + kafkaCanalConnector.ack(1); + } + } + }); + + sleep(120000); + running = false; + kafkaCanalConnector.disconnect(); + } + +} diff --git a/kafka-client/src/test/resources/logback.xml b/kafka-client/src/test/resources/logback.xml new file mode 100644 index 00000000..81fa0712 --- /dev/null +++ b/kafka-client/src/test/resources/logback.xml @@ -0,0 +1,19 @@ + + + + + + %d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{56} - %msg%n + + + + + + + + + + + + + \ No newline at end of file diff --git a/kafka/pom.xml b/kafka/pom.xml index 55e72553..68d6d72e 100644 --- a/kafka/pom.xml +++ b/kafka/pom.xml @@ -26,12 +26,6 @@ 1.17 - - - - - - org.apache.kafka kafka_2.11 diff --git a/kafka/src/main/resources/kafka.yml b/kafka/src/main/resources/kafka.yml index 7ca1bb7f..a869321c 100644 --- a/kafka/src/main/resources/kafka.yml +++ b/kafka/src/main/resources/kafka.yml @@ -5,8 +5,11 @@ lingerMs: 1 bufferMemory: 33554432 topics: - - topic: example + - topic: expTest partition: canalDestination: example +# - topic: example2 +# partition: +# canalDestination: example diff --git a/pom.xml b/pom.xml index bbe33de0..e8a97b45 100644 --- a/pom.xml +++ b/pom.xml @@ -128,6 +128,7 @@ deployer example kafka + kafka-client