From db82be9abaea413030a151feb46a22bdc3f2ea95 Mon Sep 17 00:00:00 2001 From: YunaiV <> Date: Sat, 7 Dec 2019 08:53:54 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20spring=20boot=20kafka=20?= =?UTF-8?q?=E7=A4=BA=E4=BE=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../kafkademo/consumer/Demo02Consumer.java | 2 +- .../kafkademo/config/KafkaConfiguration.java | 28 +++++++++++ .../kafkademo/consumer/Demo04Consumer.java | 42 ++++++++++++++++ .../kafkademo/message/Demo04Message.java | 31 ++++++++++++ .../kafkademo/producer/Demo04Producer.java | 25 ++++++++++ .../producer/Demo04ProducerTest.java | 48 +++++++++++++++++++ 6 files changed, 175 insertions(+), 1 deletion(-) create mode 100644 lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/config/KafkaConfiguration.java create mode 100644 lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo04Consumer.java create mode 100644 lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/message/Demo04Message.java create mode 100644 lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04Producer.java create mode 100644 lab-03/lab-03-kafka-demo/src/test/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04ProducerTest.java diff --git a/lab-03/lab-03-kafka-demo-batch-consume/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo02Consumer.java b/lab-03/lab-03-kafka-demo-batch-consume/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo02Consumer.java index d6c095d6..5a412107 100644 --- a/lab-03/lab-03-kafka-demo-batch-consume/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo02Consumer.java +++ b/lab-03/lab-03-kafka-demo-batch-consume/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo02Consumer.java @@ -14,7 +14,7 @@ public class Demo02Consumer { private Logger logger = LoggerFactory.getLogger(getClass()); @KafkaListener(topics = Demo02Message.TOPIC, - groupId = "demo02-B-consumer-group-" + Demo02Message.TOPIC) + groupId = "demo02-consumer-group-" + Demo02Message.TOPIC) public void onMessage(List messages) { logger.info("[onMessage][线程编号:{} 消息数量:{}]", Thread.currentThread().getId(), messages.size()); } diff --git a/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/config/KafkaConfiguration.java b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/config/KafkaConfiguration.java new file mode 100644 index 00000000..285a8ac7 --- /dev/null +++ b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/config/KafkaConfiguration.java @@ -0,0 +1,28 @@ +package cn.iocoder.springboot.lab03.kafkademo.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.ConsumerRecordRecoverer; +import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; +import org.springframework.kafka.listener.ErrorHandler; +import org.springframework.kafka.listener.SeekToCurrentErrorHandler; +import org.springframework.util.backoff.BackOff; +import org.springframework.util.backoff.FixedBackOff; + +@Configuration +public class KafkaConfiguration { + + @Bean + @Primary + public ErrorHandler kafkaErrorHandler(KafkaTemplate template) { + // 创建 DeadLetterPublishingRecoverer 对象 + ConsumerRecordRecoverer recoverer = new DeadLetterPublishingRecoverer(template); + // 创建 FixedBackOff 对象 + BackOff backOff = new FixedBackOff(10 * 1000L, 3L); + // 创建 SeekToCurrentErrorHandler 对象 + return new SeekToCurrentErrorHandler(recoverer, backOff); + } + +} diff --git a/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo04Consumer.java b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo04Consumer.java new file mode 100644 index 00000000..f474ca46 --- /dev/null +++ b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/consumer/Demo04Consumer.java @@ -0,0 +1,42 @@ +package cn.iocoder.springboot.lab03.kafkademo.consumer; + +import cn.iocoder.springboot.lab03.kafkademo.message.Demo04Message; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; + +import java.util.concurrent.atomic.AtomicInteger; + +@Component +public class Demo04Consumer { + + private AtomicInteger count = new AtomicInteger(0); + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @KafkaListener(topics = Demo04Message.TOPIC, + groupId = "demo04-consumer-group-" + Demo04Message.TOPIC) + public void onMessage(Demo04Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + // 注意,此处抛出一个 RuntimeException 异常,模拟消费失败 + throw new RuntimeException("我就是故意抛出一个异常"); + } + +// @Bean +// public ConsumerAwareListenerErrorHandler listenErrorHandler() { +// return new ConsumerAwareListenerErrorHandler() { +// +// @Override +// public Object handleError(Message message, +// ListenerExecutionFailedException e, +// Consumer consumer) { +// System.out.println("message:" + message.getPayload()); +// System.out.println("exception:" + e.getMessage()); +// return null; +// } +// +// }; +// } + +} diff --git a/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/message/Demo04Message.java b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/message/Demo04Message.java new file mode 100644 index 00000000..beea98c0 --- /dev/null +++ b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/message/Demo04Message.java @@ -0,0 +1,31 @@ +package cn.iocoder.springboot.lab03.kafkademo.message; + +/** + * 示例 04 的 Message 消息 + */ +public class Demo04Message { + + public static final String TOPIC = "DEMO_04"; + + /** + * 编号 + */ + private Integer id; + + public Demo04Message setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "Demo01Message{" + + "id=" + id + + '}'; + } + +} diff --git a/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04Producer.java b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04Producer.java new file mode 100644 index 00000000..78c29e8f --- /dev/null +++ b/lab-03/lab-03-kafka-demo/src/main/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04Producer.java @@ -0,0 +1,25 @@ +package cn.iocoder.springboot.lab03.kafkademo.producer; + +import cn.iocoder.springboot.lab03.kafkademo.message.Demo04Message; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.SendResult; +import org.springframework.stereotype.Component; + +import javax.annotation.Resource; +import java.util.concurrent.ExecutionException; + +@Component +public class Demo04Producer { + + @Resource + private KafkaTemplate kafkaTemplate; + + public SendResult syncSend(Integer id) throws ExecutionException, InterruptedException { + // 创建 Demo04Message 消息 + Demo04Message message = new Demo04Message(); + message.setId(id); + // 同步发送消息 + return kafkaTemplate.send(Demo04Message.TOPIC, message).get(); + } + +} diff --git a/lab-03/lab-03-kafka-demo/src/test/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04ProducerTest.java b/lab-03/lab-03-kafka-demo/src/test/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04ProducerTest.java new file mode 100644 index 00000000..8c2b44ca --- /dev/null +++ b/lab-03/lab-03-kafka-demo/src/test/java/cn/iocoder/springboot/lab03/kafkademo/producer/Demo04ProducerTest.java @@ -0,0 +1,48 @@ +package cn.iocoder.springboot.lab03.kafkademo.producer; + +import cn.iocoder.springboot.lab03.kafkademo.Application; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.kafka.support.SendResult; +import org.springframework.test.context.junit4.SpringRunner; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; + +@RunWith(SpringRunner.class) +@SpringBootTest(classes = Application.class) +public class Demo04ProducerTest { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @Autowired + private Demo04Producer producer; + + @Test + public void testSyncSend() throws ExecutionException, InterruptedException { + int id = (int) (System.currentTimeMillis() / 1000); + SendResult result = producer.syncSend(id); + logger.info("[testSyncSend][发送编号:[{}] 发送结果:[{}]]", id, result); + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + + @Test + public void testSyncSendX() throws ExecutionException, InterruptedException { + for (int i = 0; i < 100; i++) { + int id = (int) (System.currentTimeMillis() / 1000); + SendResult result = producer.syncSend(id); + logger.info("[testSyncSend][发送编号:[{}] 发送结果:[{}]]", id, result); + Thread.sleep(10 * 1000L); + } + + // 阻塞等待,保证消费 + new CountDownLatch(1).await(); + } + +}