diff --git a/lab-39/lab-39-kafka/pom.xml b/lab-39/lab-39-kafka/pom.xml new file mode 100644 index 00000000..6b0f942d --- /dev/null +++ b/lab-39/lab-39-kafka/pom.xml @@ -0,0 +1,31 @@ + + + + org.springframework.boot + spring-boot-starter-parent + 2.2.2.RELEASE + + + 4.0.0 + + lab-39-kafka + + + + + + org.springframework.kafka + spring-kafka + 2.3.3.RELEASE + + + + + org.springframework.boot + spring-boot-starter-web + + + + diff --git a/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/KafkaApplication.java b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/KafkaApplication.java new file mode 100644 index 00000000..692e130e --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/KafkaApplication.java @@ -0,0 +1,13 @@ +package cn.iocoder.springboot.lab39.skywalkingdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class KafkaApplication { + + public static void main(String[] args) { + SpringApplication.run(KafkaApplication.class, args); + } + +} diff --git a/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/consumer/DemoConsumer.java b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/consumer/DemoConsumer.java new file mode 100644 index 00000000..7c77a87b --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/consumer/DemoConsumer.java @@ -0,0 +1,20 @@ +package cn.iocoder.springboot.lab39.skywalkingdemo.consumer; + +import cn.iocoder.springboot.lab39.skywalkingdemo.message.DemoMessage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; + +@Component +public class DemoConsumer { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @KafkaListener(topics = DemoMessage.TOPIC, + groupId = "demo-consumer-group-" + DemoMessage.TOPIC) + public void onMessage(DemoMessage message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/controller/DemoController.java b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/controller/DemoController.java new file mode 100644 index 00000000..39afd3c2 --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/controller/DemoController.java @@ -0,0 +1,28 @@ +package cn.iocoder.springboot.lab39.skywalkingdemo.controller; + +import cn.iocoder.springboot.lab39.skywalkingdemo.producer.DemoProducer; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +import java.util.concurrent.ExecutionException; + +@RestController +@RequestMapping("/demo") +public class DemoController { + + @Autowired + private DemoProducer producer; + + @GetMapping("/kafka") + public String echo() throws ExecutionException, InterruptedException { + this.sendMessage(1); + return "kafka"; + } + + public void sendMessage(Integer id) throws ExecutionException, InterruptedException { + producer.syncSend(id); + } + +} diff --git a/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/message/DemoMessage.java b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/message/DemoMessage.java new file mode 100644 index 00000000..51f5065c --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/message/DemoMessage.java @@ -0,0 +1,31 @@ +package cn.iocoder.springboot.lab39.skywalkingdemo.message; + +/** + * 示例 Message 消息 + */ +public class DemoMessage { + + public static final String TOPIC = "DEMO"; + + /** + * 编号 + */ + private Integer id; + + public DemoMessage setId(Integer id) { + this.id = id; + return this; + } + + public Integer getId() { + return id; + } + + @Override + public String toString() { + return "DemoMessage{" + + "id=" + id + + '}'; + } + +} diff --git a/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/producer/DemoProducer.java b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/producer/DemoProducer.java new file mode 100644 index 00000000..072b4243 --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/java/cn/iocoder/springboot/lab39/skywalkingdemo/producer/DemoProducer.java @@ -0,0 +1,25 @@ +package cn.iocoder.springboot.lab39.skywalkingdemo.producer; + +import cn.iocoder.springboot.lab39.skywalkingdemo.message.DemoMessage; +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 DemoProducer { + + @Resource + private KafkaTemplate kafkaTemplate; + + public SendResult syncSend(Integer id) throws ExecutionException, InterruptedException { + // 创建 DemoMessage 消息 + DemoMessage message = new DemoMessage(); + message.setId(id); + // 同步发送消息 + return kafkaTemplate.send(DemoMessage.TOPIC, message).get(); + } + +} diff --git a/lab-39/lab-39-kafka/src/main/resources/application.yaml b/lab-39/lab-39-kafka/src/main/resources/application.yaml new file mode 100644 index 00000000..3abbadd4 --- /dev/null +++ b/lab-39/lab-39-kafka/src/main/resources/application.yaml @@ -0,0 +1,26 @@ +server: + port: 8079 + +spring: + # Kafka 配置项,对应 KafkaProperties 配置类 + kafka: + bootstrap-servers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔 + # Kafka Producer 配置项 + producer: + acks: 1 # 0-不应答。1-leader 应答。all-所有 leader 和 follower 应答。 + retries: 3 # 发送失败时,重试发送的次数 + key-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息的 key 的序列化 + value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 消息的 value 的序列化 + # Kafka Consumer 配置项 + consumer: + auto-offset-reset: earliest # 设置消费者分组最初的消费进度为 earliest 。可参考博客 https://blog.csdn.net/lishuangzhe7047/article/details/74530417 理解 + key-deserializer: org.apache.kafka.common.serialization.StringDeserializer + value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer + properties: + spring: + json: + trusted: + packages: cn.iocoder.springboot.lab39.skywalkingdemo.message # 消息 POJO 可信目录,解决 JSON 无法反序列化的问题 + # Kafka Consumer Listener 监听器配置 + listener: + missing-topics-fatal: false # 消费监听接口监听的主题不存在时,默认会报错。所以通过设置为 false ,解决报错 diff --git a/lab-39/lab-39-kafka/target/classes/application.yaml b/lab-39/lab-39-kafka/target/classes/application.yaml new file mode 100644 index 00000000..3abbadd4 --- /dev/null +++ b/lab-39/lab-39-kafka/target/classes/application.yaml @@ -0,0 +1,26 @@ +server: + port: 8079 + +spring: + # Kafka 配置项,对应 KafkaProperties 配置类 + kafka: + bootstrap-servers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔 + # Kafka Producer 配置项 + producer: + acks: 1 # 0-不应答。1-leader 应答。all-所有 leader 和 follower 应答。 + retries: 3 # 发送失败时,重试发送的次数 + key-serializer: org.apache.kafka.common.serialization.StringSerializer # 消息的 key 的序列化 + value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 消息的 value 的序列化 + # Kafka Consumer 配置项 + consumer: + auto-offset-reset: earliest # 设置消费者分组最初的消费进度为 earliest 。可参考博客 https://blog.csdn.net/lishuangzhe7047/article/details/74530417 理解 + key-deserializer: org.apache.kafka.common.serialization.StringDeserializer + value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer + properties: + spring: + json: + trusted: + packages: cn.iocoder.springboot.lab39.skywalkingdemo.message # 消息 POJO 可信目录,解决 JSON 无法反序列化的问题 + # Kafka Consumer Listener 监听器配置 + listener: + missing-topics-fatal: false # 消费监听接口监听的主题不存在时,默认会报错。所以通过设置为 false ,解决报错 diff --git a/lab-39/pom.xml b/lab-39/pom.xml index 25477c12..a403a343 100644 --- a/lab-39/pom.xml +++ b/lab-39/pom.xml @@ -20,6 +20,7 @@ lab-39-elasticsearch lab-39-elasticsearch-jest lab-39-rocketmq + lab-39-kafka