diff --git a/lab-03/pom.xml b/lab-03/pom.xml index c34541e6..a6aec689 100644 --- a/lab-03/pom.xml +++ b/lab-03/pom.xml @@ -3,15 +3,45 @@ xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> - labs-parent - cn.iocoder.springboot.labs - 1.0-SNAPSHOT + org.springframework.boot + spring-boot-starter-parent + 1.5.9.RELEASE + 4.0.0 lab-03 + + + + org.apache.maven.plugins + maven-compiler-plugin + + 8 + 8 + + + + + + org.springframework.boot + spring-boot-starter-web + + + + org.projectlombok + lombok + true + + + + org.springframework.boot + spring-boot-starter-test + test + + org.apache.kafka kafka-clients @@ -22,6 +52,36 @@ fastjson 1.2.55 + + org.springframework.integration + spring-integration-kafka + 3.1.0.RELEASE + + + + org.springframework.kafka + spring-kafka + 1.1.1.RELEASE + + + + com.google.code.gson + gson + 2.8.2 + + + + org.projectlombok + lombok + true + + + + org.springframework.boot + spring-boot-starter-test + test + + \ No newline at end of file diff --git a/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaApplication.java b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaApplication.java new file mode 100644 index 00000000..dab9c905 --- /dev/null +++ b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaApplication.java @@ -0,0 +1,34 @@ +package cn.iocoder.springboot.labs.lab03; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.ConfigurableApplicationContext; + +@SpringBootApplication +public class KafkaApplication { + + public static void main(String[] args) { + + ConfigurableApplicationContext context = SpringApplication.run(KafkaApplication.class, args); + +// KafkaSender sender = context.getBean(KafkaSender.class); +// +// for (int i = 0; i < 3; i++) { +// //调用消息发送类中的消息发送方法 +// sender.send(); +// +// try { +// Thread.sleep(3000); +// } catch (InterruptedException e) { +// e.printStackTrace(); +// } +// } + + try { + Thread.sleep(300000); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + +} diff --git a/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaReceiver.java b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaReceiver.java new file mode 100644 index 00000000..3eae5e71 --- /dev/null +++ b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/KafkaReceiver.java @@ -0,0 +1,61 @@ +package cn.iocoder.springboot.labs.lab03; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; + +import java.util.Optional; + +@Component +public class KafkaReceiver { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @KafkaListener(topics = {"test"}) + public void listen(ConsumerRecord record) { + + Optional kafkaMessage = Optional.ofNullable(record.value()); + + if (kafkaMessage.isPresent()) { + + Object message = kafkaMessage.get(); + + logger.info("[test{}]----------------- record =" + record, Thread.currentThread().getId()); + logger.info("[test{}]------------------ message =" + message, Thread.currentThread().getId()); + } + + } + + @KafkaListener(topics = {"yunai"}) + public void listen02(ConsumerRecord record) { + + Optional kafkaMessage = Optional.ofNullable(record.value()); + + if (kafkaMessage.isPresent()) { + + Object message = kafkaMessage.get(); + + logger.info("[yunai{}]----------------- record =" + record, Thread.currentThread().getId()); + logger.info("[yunai{}]------------------ message =" + message, Thread.currentThread().getId()); + } + + } + + @KafkaListener(topics = {"afei"}) + public void listen03(ConsumerRecord record) { + + Optional kafkaMessage = Optional.ofNullable(record.value()); + + if (kafkaMessage.isPresent()) { + + Object message = kafkaMessage.get(); + + logger.info("[afei{}]----------------- record =" + record, Thread.currentThread().getId()); + logger.info("[afei{}]------------------ message =" + message, Thread.currentThread().getId()); + } + + } + +} diff --git a/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/pojo/Message.java b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/pojo/Message.java new file mode 100644 index 00000000..0c24863c --- /dev/null +++ b/lab-03/src/main/java/cn/iocoder/springboot/labs/lab03/pojo/Message.java @@ -0,0 +1,16 @@ +package cn.iocoder.springboot.labs.lab03.pojo; + +import lombok.Data; + +import java.util.Date; + +@Data +public class Message { + + private Long id; //id + + private String msg; //消息 + + private Date sendTime; //时间戳 + +} diff --git a/lab-03/src/main/resources/application.properties b/lab-03/src/main/resources/application.properties new file mode 100644 index 00000000..98ecb9ed --- /dev/null +++ b/lab-03/src/main/resources/application.properties @@ -0,0 +1,29 @@ +#============== kafka =================== +# 指定kafka 代理地址,可以多个 +spring.kafka.bootstrap-servers=127.0.0.1:9092 + +#=============== provider ======================= + +spring.kafka.producer.retries=0 +# 每次批量发送消息的数量 +spring.kafka.producer.batch-size=16384 +spring.kafka.producer.buffer-memory=33554432 + +# 指定消息key和消息体的编解码方式 +spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer +spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer + +#=============== consumer ======================= +# 指定默认消费者group id +spring.kafka.consumer.group-id=test-consumer-group6 + +spring.kafka.listener.concurrency=10 + +spring.kafka.consumer.auto-offset-reset=earliest +spring.kafka.consumer.enable-auto-commit=true +spring.kafka.consumer.auto-commit-interval=100 +enable.auto.commit=false + +# 指定消息key和消息体的编解码方式 +spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer +spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer