diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/pom.xml b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/pom.xml
new file mode 100644
index 00000000..3eea5acd
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/pom.xml
@@ -0,0 +1,58 @@
+
+
+
+ labx-11
+ cn.iocoder.springboot.labs
+ 1.0-SNAPSHOT
+
+ 4.0.0
+
+ labx-11-sc-stream-kafka-consumer-concurrency
+
+
+ 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-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java
new file mode 100644
index 00000000..dc6118d5
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java
@@ -0,0 +1,16 @@
+package cn.iocoder.springcloud.labx11.kafkademo.consumerdemo;
+
+import cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.listener.MySink;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.cloud.stream.annotation.EnableBinding;
+
+@SpringBootApplication
+@EnableBinding(MySink.class)
+public class ConsumerApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(ConsumerApplication.class, args);
+ }
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java
new file mode 100644
index 00000000..979fb543
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java
@@ -0,0 +1,20 @@
+package cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.listener;
+
+import cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.message.Demo01Message;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.cloud.stream.annotation.StreamListener;
+import org.springframework.messaging.handler.annotation.Payload;
+import org.springframework.stereotype.Component;
+
+@Component
+public class Demo01Consumer {
+
+ private Logger logger = LoggerFactory.getLogger(getClass());
+
+ @StreamListener(MySink.DEMO01_INPUT)
+ public void onMessage(@Payload Demo01Message message) {
+ logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message);
+ }
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java
new file mode 100644
index 00000000..1e510b79
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java
@@ -0,0 +1,13 @@
+package cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.listener;
+
+import org.springframework.cloud.stream.annotation.Input;
+import org.springframework.messaging.SubscribableChannel;
+
+public interface MySink {
+
+ String DEMO01_INPUT = "demo01-input";
+
+ @Input(DEMO01_INPUT)
+ SubscribableChannel demo01Input();
+
+}
diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java
new file mode 100644
index 00000000..cf57407a
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java
@@ -0,0 +1,29 @@
+package cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.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-consumer-concurrency/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/resources/application.yml
new file mode 100644
index 00000000..80cc70e9
--- /dev/null
+++ b/labx-11/labx-11-sc-stream-kafka-consumer-concurrency/src/main/resources/application.yml
@@ -0,0 +1,25 @@
+spring:
+ application:
+ name: demo-consumer-application
+ cloud:
+ # Spring Cloud Stream 配置项,对应 BindingServiceProperties 类
+ stream:
+ # Binder 配置项,对应 BinderProperties Map
+# binders:
+ # Binding 配置项,对应 BindingProperties Map
+ bindings:
+ demo01-input:
+ destination: DEMO-TOPIC-01 # 目的地。这里使用 Kafka Topic
+ content-type: application/json # 内容格式。这里使用 JSON
+ group: demo01-consumer-group # 消费者分组
+ # Consumer 配置项,对应 ConsumerProperties 类
+ consumer:
+ concurrency: 2 # 每个 Consumer 消费线程数的初始大小,默认为 1
+ # Spring Cloud Stream Kafka 配置项
+ kafka:
+ # Kafka Binder 配置项,对应 KafkaBinderConfigurationProperties 类
+ binder:
+ brokers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔
+
+server:
+ port: ${random.int[10000,19999]} # 随机端口,方便启动多个消费者
diff --git a/labx-11/pom.xml b/labx-11/pom.xml
index 92cd7944..fd418177 100644
--- a/labx-11/pom.xml
+++ b/labx-11/pom.xml
@@ -21,7 +21,11 @@
labx-11-sc-stream-kafka-consumer-error-handler
+
labx-11-sc-stream-kafka-consumer-broadcasting
+
+
+ labx-11-sc-stream-kafka-consumer-concurrency