From 7b6e6ecf080ce5b548b85d435c089cff609b04a3 Mon Sep 17 00:00:00 2001
From: YunaiV <>
Date: Tue, 10 Mar 2020 09:14:04 +0800
Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20spring=20cloud=20stream=20?=
=?UTF-8?q?kafka=20=E7=A4=BA=E4=BE=8B?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../pom.xml | 58 +++++++++++++++++++
.../kafkademo/ProducerApplication.java | 16 +++++
.../controller/Demo01Controller.java | 42 ++++++++++++++
.../kafkademo/message/Demo01Message.java | 29 ++++++++++
.../kafkademo/kafkademo/message/MySource.java | 11 ++++
.../src/main/resources/application.yml | 28 +++++++++
labx-11/pom.xml | 2 +
7 files changed, 186 insertions(+)
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/pom.xml
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java
create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/resources/application.yml
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/pom.xml b/labx-11/labx-11-sc-stream-kafka-producer-batch/pom.xml
new file mode 100644
index 00000000..579473ee
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/pom.xml
@@ -0,0 +1,58 @@
+
+
+
+ labx-11
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-11-sc-stream-kafka-producer-batch
+
+
+ 1.8
+ 1.8
+ 2.2.4.RELEASE
+ Hoxton.SR1
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ ${spring.boot.version}
+ pom
+ import
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring.cloud.version}
+ pom
+ import
+
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+
+ org.springframework.cloud
+ spring-cloud-starter-stream-kafka
+
+
+
+
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java
new file mode 100644
index 00000000..2dad6085
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloud.labx11.kafkademo.kafkademo;
+
+import cn.iocoder.springcloud.labx11.kafkademo.kafkademo.message.MySource;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.cloud.stream.annotation.EnableBinding;
+
+@SpringBootApplication
+@EnableBinding(MySource.class)
+public class ProducerApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(ProducerApplication.class, args);
+ }
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java
new file mode 100644
index 00000000..42f9d081
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java
@@ -0,0 +1,42 @@
+package cn.iocoder.springcloud.labx11.kafkademo.kafkademo.controller;
+
+import cn.iocoder.springcloud.labx11.kafkademo.kafkademo.message.Demo01Message;
+import cn.iocoder.springcloud.labx11.kafkademo.kafkademo.message.MySource;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.Random;
+
+@RestController
+@RequestMapping("/demo01")
+public class Demo01Controller {
+
+ private Logger logger = LoggerFactory.getLogger(getClass());
+
+ @Autowired
+ private MySource mySource;
+
+ @GetMapping("/send_batch")
+ public boolean sendTag() {
+ for (int i = 0; i < 3; i++) {
+ // 创建 Message
+ int id = new Random().nextInt();
+ Demo01Message message = new Demo01Message()
+ .setId(id);
+ // 创建 Spring Message 对象
+ Message springMessage = MessageBuilder.withPayload(message)
+ .build();
+ // 发送消息
+ mySource.demo01Output().send(springMessage);
+ logger.info("[send_transaction][发送编号:[{}] 发送成功]", id);
+ }
+ return true;
+ }
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java
new file mode 100644
index 00000000..6efa4249
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloud.labx11.kafkademo.kafkademo.message;
+
+/**
+ * 示例 01 的 Message 消息
+ */
+public class Demo01Message {
+
+ /**
+ * 编号
+ */
+ private Integer id;
+
+ public Demo01Message setId(Integer id) {
+ this.id = id;
+ return this;
+ }
+
+ public Integer getId() {
+ return id;
+ }
+
+ @Override
+ public String toString() {
+ return "Demo01Message{" +
+ "id=" + id +
+ '}';
+ }
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java
new file mode 100644
index 00000000..c1d58248
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java
@@ -0,0 +1,11 @@
+package cn.iocoder.springcloud.labx11.kafkademo.kafkademo.message;
+
+import org.springframework.cloud.stream.annotation.Output;
+import org.springframework.messaging.MessageChannel;
+
+public interface MySource {
+
+ @Output("demo01-output")
+ MessageChannel demo01Output();
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/resources/application.yml
new file mode 100644
index 00000000..6075ffe3
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-producer-batch/src/main/resources/application.yml
@@ -0,0 +1,28 @@
+spring:
+ application:
+ name: demo-producer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binder 配置项,对应 BinderProperties Map
+# binders:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-output:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 Kafka Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ # Spring Cloud Stream Kafka 配置项
+ kafka:
+ # Kafka Binder 配置项,对应 KafkaBinderConfigurationProperties 类
+ binder:
+ brokers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔
+ # Kafka 自定义 Binding 配置项,对应 KafkaBindingProperties Map
+ bindings:
+ demo01-output:
+ # Kafka Producer 配置项,对应 KafkaProducerProperties 类
+ producer:
+ batch-timeout: 30000 # 批处理延迟时间上限。这里配置为 30 * 1000 ms 过后,不管是否消息数量是否到达 batch-size 或者消息大小到达 buffer-memory 后,都直接发送一次请求
+ buffer-size: 33554432 # 每次批量发送消息的最大内存
+
+server:
+ port: 18080
diff --git a/labx-11/pom.xml b/labx-11/pom.xml
index 6308fb0f..1c6e41eb 100644
--- a/labx-11/pom.xml
+++ b/labx-11/pom.xml
@@ -39,6 +39,8 @@
labx-11-sc-stream-kafka-consumer-ack
+
+ labx-11-sc-stream-kafka-producer-batch