From 33a58e8aacce6335ea344064983466cde1f3ad2d Mon Sep 17 00:00:00 2001 From: hz21056617 <1IB76X4g7T> Date: Fri, 15 Apr 2022 11:08:06 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0kafka=20=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../queue/BlockingQueueTestWork.java | 2 + kafka/pom.xml | 23 ++ .../com/wyl/kafka/KafkaConfigConstant.java | 51 ++++ .../kafka/consumer/KafkaConsumerExample.java | 50 ++++ .../wyl/kafka/consumer/KafkaConsumerPoll.java | 222 ++++++++++++++++++ .../kafka/producers/KafkaProducerExample.java | 55 +++++ .../kafka/producers/KafkaProducerSend.java | 79 +++++++ kafka/src/main/resources/logback.xml | 59 +++++ .../java/com/wyl/kafka/test/ConsumerTest.java | 88 +++++++ .../java/com/wyl/kafka/test/ProducerTest.java | 24 ++ pom.xml | 18 ++ 11 files changed, 671 insertions(+) create mode 100644 kafka/pom.xml create mode 100644 kafka/src/main/java/com/wyl/kafka/KafkaConfigConstant.java create mode 100644 kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerExample.java create mode 100644 kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerPoll.java create mode 100644 kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerExample.java create mode 100644 kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerSend.java create mode 100644 kafka/src/main/resources/logback.xml create mode 100644 kafka/src/test/java/com/wyl/kafka/test/ConsumerTest.java create mode 100644 kafka/src/test/java/com/wyl/kafka/test/ProducerTest.java diff --git a/concurrent/src/main/java/concurrent/queue/BlockingQueueTestWork.java b/concurrent/src/main/java/concurrent/queue/BlockingQueueTestWork.java index e7f6e56..92f5be4 100644 --- a/concurrent/src/main/java/concurrent/queue/BlockingQueueTestWork.java +++ b/concurrent/src/main/java/concurrent/queue/BlockingQueueTestWork.java @@ -4,6 +4,7 @@ import lombok.AllArgsConstructor; import lombok.SneakyThrows; import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.ConcurrentLinkedQueue; @AllArgsConstructor public class BlockingQueueTestWork extends Thread { @@ -27,6 +28,7 @@ public class BlockingQueueTestWork extends Thread { } public void offerTest() { + ConcurrentLinkedQueue concurrentLinkedQueue = new ConcurrentLinkedQueue<>(); boolean offer = arrayBlockingQueue.offer(1); if (!offer) { System.out.println("队列满了"); diff --git a/kafka/pom.xml b/kafka/pom.xml new file mode 100644 index 0000000..dc0057d --- /dev/null +++ b/kafka/pom.xml @@ -0,0 +1,23 @@ + + + + JavaBasiceDemo + com.wyl.example + 1.0-SNAPSHOT + + 4.0.0 + kafka + + + 8 + 8 + + + + org.apache.kafka + kafka-clients + + + \ No newline at end of file diff --git a/kafka/src/main/java/com/wyl/kafka/KafkaConfigConstant.java b/kafka/src/main/java/com/wyl/kafka/KafkaConfigConstant.java new file mode 100644 index 0000000..5ceb452 --- /dev/null +++ b/kafka/src/main/java/com/wyl/kafka/KafkaConfigConstant.java @@ -0,0 +1,51 @@ +package com.wyl.kafka; + +/** + * kafka配置 + * @ClassName: KafkaConfig + * @Date: 2022/4/12 14:11 + * @author wangyl + * @version V1.0 + */ +public class KafkaConfigConstant { + /** + *kafka的broker地址 + */ + public static final String BOOTSTRAP_SERVERS = "localhost:9092"; + /** + *kafka的topic + */ + public static final String TOPIC = "test"; + /** + *kafka的groupId + */ + public static final String GROUP_ID = "test-group"; + /** + *kafka的key的序列化方式 + */ + public static final String KEY_DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer"; + /** + *kafka的value的序列化方式 + */ + public static final String VALUE_DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer"; + /** + * ACK策略 + */ + public static final String ACKS = "all"; + /** + *kafka的消息发送最大重试次数 + */ + public static final String RETRIES = "3"; + /** + *kafka的消息发送最大延时 + */ + public static final String BATCH_SIZE = "16384"; + /** + *kafka的消息发送最大延时 + */ + public static final String LINGER_MS = "1"; + public static final String BUFFER_MEMORY = "33554432"; + public static final String AUTO_COMMIT_INTERVAL = "1000"; + public static final String AUTO_OFFSET_RESET = "earliest"; + public static final String ENABLE_AUTO_COMMIT = "true"; +} diff --git a/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerExample.java b/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerExample.java new file mode 100644 index 0000000..90be345 --- /dev/null +++ b/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerExample.java @@ -0,0 +1,50 @@ +package com.wyl.kafka.consumer; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerConfig; + +import java.util.Properties; + +/** + * kafak生产者示例 + * @ClassName: KafkaProducer + * @Date: 2022/4/12 14:36 + * @author wangyl + * @version V1.0 + */ +public class KafkaConsumerExample { + + /** + * 通用消费者 + */ + public static Consumer KafkaConsumerNormal() { + Properties props = NormalProperties(); + Consumer consumerNormal = new KafkaConsumer<>(props); + return consumerNormal; + } + /** + * 手动提交消费者 + */ + public static Consumer KafkaConsumerManualCommit() { + Properties props = NormalProperties(); + props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + Consumer consumerNormal = new KafkaConsumer<>(props); + return consumerNormal; + } + + /** + * 通用配置 + */ + private static Properties NormalProperties() { + Properties props = new Properties(); + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "172.20.0.187:9092,172.20.0.188:9092,172.20.0.189:9092"); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); + props.put(ConsumerConfig.GROUP_ID_CONFIG, "wylll"); + return props; + } +} diff --git a/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerPoll.java b/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerPoll.java new file mode 100644 index 0000000..1e2f6c9 --- /dev/null +++ b/kafka/src/main/java/com/wyl/kafka/consumer/KafkaConsumerPoll.java @@ -0,0 +1,222 @@ +package com.wyl.kafka.consumer; + +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; +import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.TopicPartition; + +import java.time.Duration; +import java.util.*; +import java.util.stream.Stream; + +/** + * kafka发送消息测试 + * @ClassName: KafkaProducerSend + * @Date: 2022/4/12 15:18 + * @author wangyl + * @version V1.0 + */ +@Slf4j +public class KafkaConsumerPoll { + /** + * 消费者消费消息 + * @param kafkaConsumerNormal + * @param topic + * @return void + * @Date 2022/4/12 17:43 + * @Author wangyl + * @Version V1.0 + */ + public static void PollMessage(Consumer kafkaConsumerNormal, String topic) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic)); + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + } + } + + /** + * 同步手动提交 + * @param kafkaConsumerNormal + * @param topic + * @return void + * @Date 2022/4/13 18:44 + * @Author wangyl + * @Version V1.0 + */ + public static void PollMessageCommitSync(Consumer kafkaConsumerNormal, String topic) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic)); + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + kafkaConsumerNormal.commitSync(); + } + + } + + /** + * 同步手动提交offset + * @param kafkaConsumerNormal + * @param topic + * @return void + * @Date 2022/4/13 18:44 + * @Author wangyl + * @Version V1.0 + */ + public static void PollMessageCommitSyncOffsets(Consumer kafkaConsumerNormal, String topic) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic)); + Map offsets = new HashMap<>(); + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)); + }); + kafkaConsumerNormal.commitSync(offsets); + offsets.clear(); + } + } + + /** + * 异步提交,有回调 + * @param kafkaConsumerNormal + * @param topic + * @return void + * @Date 2022/4/13 18:47 + * @Author wangyl + * @Version V1.0 + */ + public static void PollMessageCommitAsyncCallBack(Consumer kafkaConsumerNormal, String topic) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic)); + Map offsets = new HashMap<>(); + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)); + }); + kafkaConsumerNormal.commitAsync(offsets, (offsetsInfo, exception) -> { + if (exception != null) { + log.error("提交消息异常:{}", exception.getMessage()); + } else { + log.info("提交消息成功:{}", offsetsInfo); + } + }); + } + } + + /** + * 异步回调offsets + * @param kafkaConsumerNormal + * @param topic + * @return void + * @Date 2022/4/13 18:50 + * @Author wangyl + * @Version V1.0 + */ + public static void PollMessageCommitAsyncOffsetsCallBack(Consumer kafkaConsumerNormal, String topic) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic)); + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + kafkaConsumerNormal.commitAsync(); + } + } + + + public static void PollMessageSeek(Consumer kafkaConsumerNormal, String topic,int partition,int offset) { + //消费者消费消息 + kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() { + @Override + public void onPartitionsRevoked(Collection partitions) { + + } + + @Override + public void onPartitionsAssigned(Collection partitions) { + kafkaConsumerNormal.seek(new TopicPartition(topic, partition), offset); + } + }); + + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + kafkaConsumerNormal.commitAsync(); + } + } + + public static void PollMessageSeekBegin(Consumer kafkaConsumerNormal, String topic,int partition) { + final TopicPartition topicPartition = new TopicPartition(topic, partition); + final List topicPartitions = Arrays.asList(topicPartition); + kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() { + @Override + public void onPartitionsRevoked(Collection partitions) { + + } + + @Override + public void onPartitionsAssigned(Collection partitions) { + kafkaConsumerNormal.seekToBeginning(topicPartitions); + } + }); + + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + kafkaConsumerNormal.commitAsync(); + } + } + + public static void PollMessageSeekEnd(Consumer kafkaConsumerNormal, String topic,int partition) { + //消费者消费消息 + final TopicPartition topicPartition = new TopicPartition(topic, partition); + final List topicPartitions = Arrays.asList(topicPartition); + kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() { + @Override + public void onPartitionsRevoked(Collection partitions) { + + } + + @Override + public void onPartitionsAssigned(Collection partitions) { + kafkaConsumerNormal.seekToEnd(topicPartitions); + } + }); + + while (true) { + //消费消息 + kafkaConsumerNormal.poll(Duration.ofMillis(100)) + .forEach(record -> { + log.info("消费者消费消息:{},分区:{},topic:{},offset:{}", record.value(), record.partition(), record.topic(), record.offset()); + }); + kafkaConsumerNormal.commitAsync(); + } + } + + +} diff --git a/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerExample.java b/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerExample.java new file mode 100644 index 0000000..53bda7d --- /dev/null +++ b/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerExample.java @@ -0,0 +1,55 @@ +package com.wyl.kafka.producers; + +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerConfig; + +import java.util.Properties; + +/** + * kafak生产者示例 + * @ClassName: KafkaProducer + * @Date: 2022/4/12 14:36 + * @author wangyl + * @version V1.0 + */ +public class KafkaProducerExample { + + /** + * 通用生产者 + */ + public static Producer KafkaProducerNormal() { + Properties props = NormalProperties(); + Producer producerNormal = new KafkaProducer(props); + return producerNormal; + } + + /** + * 事务消息生产者 + */ + public static Producer KafkaProducerTransactional() { + Properties props = NormalProperties(); + /** + *设置事务 id(必须),事务 id 任意起名 + */ + props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "transactionalId"); + props.put(ProducerConfig.RETRIES_CONFIG, 2); + Producer producerNormal = new KafkaProducer(props); + return producerNormal; + } + + /** + * 通用配置 + */ + private static Properties NormalProperties() { + Properties props = new Properties(); + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "172.20.0.187:9092,172.20.0.188:9092,172.20.0.189:9092"); + props.put(ProducerConfig.ACKS_CONFIG, "all"); + props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); + props.put(ProducerConfig.LINGER_MS_CONFIG, 1); + props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); + props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); + props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); + return props; + } +} diff --git a/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerSend.java b/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerSend.java new file mode 100644 index 0000000..922d123 --- /dev/null +++ b/kafka/src/main/java/com/wyl/kafka/producers/KafkaProducerSend.java @@ -0,0 +1,79 @@ +package com.wyl.kafka.producers; + +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; + +/** + * kafka发送消息测试 + * @ClassName: KafkaProducerSend + * @Date: 2022/4/12 15:18 + * @author wangyl + * @version V1.0 + */ +@Slf4j +public class KafkaProducerSend { + /** + * 生产者发送回调消息 + * @param kafkaProducerNormal + * @param topic + * @param message + * @return void + * @Date 2022/4/12 17:43 + * @Author wangyl + * @Version V1.0 + */ + public static void sendMessageAndCallback(Producer kafkaProducerNormal, String topic, String message) { + ProducerRecord messageObj = new ProducerRecord<>(topic, message); + kafkaProducerNormal.send(messageObj, (metadata, exception) -> { + if (null != exception) { + log.error("kafka 信息推送失败[" + messageObj + "]", exception); + } + log.info("消息发送成功[{}]", messageObj); + }); + kafkaProducerNormal.close(); + } + + /** + * 发送事务消息 + * @param kafkaProducerTransactional + * @param topic + * @param message + * @return void + * @Date 2022/4/12 17:42 + * @Author wangyl + * @Version V1.0 + */ + public static void sendMessageTransactional(Producer kafkaProducerTransactional, String topic, String message) { + /** + * 初始化事务 + */ + kafkaProducerTransactional.initTransactions(); + /** + * 开启事务 + */ + kafkaProducerTransactional.beginTransaction(); + log.info("开始发送事务消息"); + try { + for (int i = 0; i < 10; i++){ + ProducerRecord messageObj = new ProducerRecord<>(topic, message + i); + kafkaProducerTransactional.send(messageObj, (metadata, exception) -> { + if (null != exception) { + log.error("kafka 信息推送失败[" + messageObj + "]", exception); + } + }); + } + log.info("消息发送成功"); + kafkaProducerTransactional.commitTransaction(); + } + catch (Exception e) { + log.error("kafka 事务消息推送失败", e); + kafkaProducerTransactional.abortTransaction(); + } + finally { + kafkaProducerTransactional.close(); + } + + } + +} diff --git a/kafka/src/main/resources/logback.xml b/kafka/src/main/resources/logback.xml new file mode 100644 index 0000000..6c0daa1 --- /dev/null +++ b/kafka/src/main/resources/logback.xml @@ -0,0 +1,59 @@ + + + + + + + + + %date{yyyy-MM-dd HH:mm:ss.SSS} java %-5level [%replace(%thread){'\s', '-'}] %logger{36} %X{trace.id} %msg%n + + utf8 + + + + + + true + + + + ${logback.output}/server.log + + ${logback.output}/server.%d{yyyy-MM-dd}.%i.log + + 1GB + 10 + 10GB + + + %date{HH:mm:ss.SSS} [%replace(%thread){'\s', '-'}] %-5level %logger{36} %X{trace.id:-huizTraceId} %msg%n + + + + + + + true + + + + + + + + + + + + + + + + + + + + + + diff --git a/kafka/src/test/java/com/wyl/kafka/test/ConsumerTest.java b/kafka/src/test/java/com/wyl/kafka/test/ConsumerTest.java new file mode 100644 index 0000000..8fc9450 --- /dev/null +++ b/kafka/src/test/java/com/wyl/kafka/test/ConsumerTest.java @@ -0,0 +1,88 @@ +package com.wyl.kafka.test; + +import com.wyl.kafka.consumer.KafkaConsumerExample; +import com.wyl.kafka.consumer.KafkaConsumerPoll; +import com.wyl.kafka.producers.KafkaProducerExample; +import com.wyl.kafka.producers.KafkaProducerSend; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.junit.jupiter.api.Test; + +@Slf4j +public class ConsumerTest { + final static String TOPIC_NAME = "wyl-filebeat-test"; + + /** + * 消费者模式 + */ + @Test + public void normalProducerSendMessageAndCallbackTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerNormal(); + KafkaConsumerPoll.PollMessage(consumer, TOPIC_NAME); + } + + /** + * 手动同步提交offset + */ + @Test + public void pollMessageCommitSyncTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerManualCommit(); + KafkaConsumerPoll.PollMessageCommitSync(consumer, TOPIC_NAME); + } + + /** + * 手动异步提交offset + */ + @Test + public void pollMessageCommitSyncOffsets() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerManualCommit(); + KafkaConsumerPoll.PollMessageCommitSyncOffsets(consumer, TOPIC_NAME); + } + + /** + * 手动异步提交offset回调 + */ + @Test + public void pollMessageCommitAsyncCallBackTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerManualCommit(); + KafkaConsumerPoll.PollMessageCommitAsyncCallBack(consumer, TOPIC_NAME); + } + + /** + * 手动异步提交offset批量回调 + */ + @Test + public void pollMessageCommitAsyncOffsetsCallBackTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerManualCommit(); + KafkaConsumerPoll.PollMessageCommitAsyncOffsetsCallBack(consumer, TOPIC_NAME); + } + + /** + * 移动到指定问位置消费消息 + */ + @Test + public void pollMessageSeekTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerNormal(); + KafkaConsumerPoll.PollMessageSeek(consumer, TOPIC_NAME, 1, 34236); + } + + /** + * 移动到开始位置消费消息 + */ + @Test + public void pollMessageSeekBeginTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerNormal(); + KafkaConsumerPoll.PollMessageSeekBegin(consumer, TOPIC_NAME, 1); + } + + /** + * 移动到结束位置消费消息 + */ + @Test + public void pollMessageCommitSeekEndTest() { + final Consumer consumer = KafkaConsumerExample.KafkaConsumerNormal(); + KafkaConsumerPoll.PollMessageSeekEnd(consumer, TOPIC_NAME, 1); + } + +} diff --git a/kafka/src/test/java/com/wyl/kafka/test/ProducerTest.java b/kafka/src/test/java/com/wyl/kafka/test/ProducerTest.java new file mode 100644 index 0000000..5afaf25 --- /dev/null +++ b/kafka/src/test/java/com/wyl/kafka/test/ProducerTest.java @@ -0,0 +1,24 @@ +package com.wyl.kafka.test; + +import com.wyl.kafka.producers.KafkaProducerExample; +import com.wyl.kafka.producers.KafkaProducerSend; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.producer.Producer; +import org.junit.jupiter.api.Test; + +@Slf4j +public class ProducerTest { + final static String TOPIC_NAME = "wyl-filebeat-test"; + + @Test + public void normalProducerSendMessageAndCallbackTest() { + final Producer stringStringProducer = KafkaProducerExample.KafkaProducerNormal(); + KafkaProducerSend.sendMessageAndCallback(stringStringProducer, TOPIC_NAME, "test123"); + } + + @Test + public void transactionalProducerSendMessageAndCallbackTest() { + final Producer stringStringProducer = KafkaProducerExample.KafkaProducerTransactional(); + KafkaProducerSend.sendMessageTransactional(stringStringProducer, TOPIC_NAME, "test123"); + } +} diff --git a/pom.xml b/pom.xml index 9ec444a..a691a25 100644 --- a/pom.xml +++ b/pom.xml @@ -12,6 +12,7 @@ elasticsearch concurrent file-read + kafka pom @@ -25,6 +26,12 @@ cglib 3.3.0 + + + org.apache.kafka + kafka-clients + 2.8.1 + @@ -39,5 +46,16 @@ 5.8.2 test + + ch.qos.logback + logback-core + 1.2.3 + + + ch.qos.logback + logback-classic + 1.2.3 + + \ No newline at end of file