From cf7c2ff10915eb31eb049a168f1898ab08f2da8d Mon Sep 17 00:00:00 2001 From: YunaiV <> Date: Mon, 9 Mar 2020 19:28:31 +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 --- .../controller/Demo01Controller.java | 16 +++++ .../controller/Demo01Controller.java | 3 +- .../src/main/resources/application.yml | 7 ++- .../pom.xml | 58 +++++++++++++++++++ .../consumerdemo/ConsumerApplication.java | 16 +++++ .../consumerdemo/listener/Demo01Consumer.java | 20 +++++++ .../consumerdemo/listener/MySink.java | 13 +++++ .../consumerdemo/message/Demo01Message.java | 29 ++++++++++ .../src/main/resources/application.yml | 22 +++++++ .../pom.xml | 58 +++++++++++++++++++ .../consumerdemo/ConsumerApplication.java | 16 +++++ .../consumerdemo/listener/Demo01Consumer.java | 19 ++++++ .../consumerdemo/listener/MySink.java | 13 +++++ .../consumerdemo/message/Demo01Message.java | 29 ++++++++++ .../src/main/resources/application.yml | 25 ++++++++ .../controller/Demo01Controller.java | 16 +++++ .../pom.xml | 58 +++++++++++++++++++ .../kafkademo/ProducerApplication.java | 16 +++++ .../controller/Demo01Controller.java | 41 +++++++++++++ .../kafkademo/message/Demo01Message.java | 29 ++++++++++ .../kafkademo/kafkademo/message/MySource.java | 11 ++++ .../src/main/resources/application.yml | 30 ++++++++++ labx-11/pom.xml | 7 +++ 23 files changed, 547 insertions(+), 5 deletions(-) create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/pom.xml create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/resources/application.yml create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/pom.xml create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java create mode 100644 labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/resources/application.yml create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/pom.xml create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java create mode 100644 labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/resources/application.yml diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java index cfee040e..084dba7b 100644 --- a/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-demo/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java @@ -34,4 +34,20 @@ public class Demo01Controller { return mySource.demo01Output().send(springMessage); } + @GetMapping("/send_tag") + public boolean sendTag() { + for (String tag : new String[]{"yunai", "yutou", "tudou"}) { + // 创建 Message + Demo01Message message = new Demo01Message() + .setId(new Random().nextInt()); + // 创建 Spring Message 对象 + Message springMessage = MessageBuilder.withPayload(message) + .setHeader("tag", tag) // 设置 Tag + .build(); + // 发送消息 + mySource.demo01Output().send(springMessage); + } + return true; + } + } diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java b/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java index 8e9f7136..89d23a23 100644 --- a/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx10/rabbitmqdemo/producerdemo/controller/Demo01Controller.java @@ -22,14 +22,13 @@ public class Demo01Controller { @Autowired private MySource mySource; - @GetMapping("/send_partition") + @GetMapping("/send_orderly") public boolean send() { // 创建 Message Demo01Message message = new Demo01Message() .setId(new Random().nextInt()); // 创建 Spring Message 对象 Message springMessage = MessageBuilder.withPayload(message) - .setHeader("partitionKey", 1) .build(); // 发送消息 return mySource.demo01Output().send(springMessage); diff --git a/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/resources/application.yml b/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/resources/application.yml index d2ab0e15..153a4479 100644 --- a/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/resources/application.yml +++ b/labx-10/labx-10-sc-stream-rabbitmq-producer-partitioning/src/main/resources/application.yml @@ -22,10 +22,11 @@ spring: destination: DEMO-TOPIC-03 # 目的地。这里使用 RabbitMQ Exchange content-type: application/json # 内容格式。这里使用 JSON binder: rabbit001 # 设置使用的 Binder 名字 + # Producer 配置项,对应 ProducerProperties 类 producer: - partition-key-expression: headers['partitionKey'] - partition-count: 4 - required-groups: demo01-consumer-group-DEMO-TOPIC-01 #only applicable for rabbit + partition-key-expression: payload['id'] # 分区 key 表达式。该表达式基于 Spring EL,从消息中获得分区 key。 + partition-count: 2 # 分区大小,默认为 1 分区 +# required-groups: demo01-consumer-group-DEMO-TOPIC-01 #only applicable for rabbit server: port: 18080 diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-filter/pom.xml b/labx-11/labx-11-sc-stream-kafka-consumer-filter/pom.xml new file mode 100644 index 00000000..1533110a --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-filter/pom.xml @@ -0,0 +1,58 @@ + + + + labx-11 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-11-sc-stream-kafka-consumer-filter + + + 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-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java b/labx-11/labx-11-sc-stream-kafka-consumer-filter/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-filter/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-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java b/labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java new file mode 100644 index 00000000..a4a0d5be --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-filter/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(value = MySink.DEMO01_INPUT, condition = "headers['tag'] == 'yunai'") + public void onMessage(@Payload Demo01Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java b/labx-11/labx-11-sc-stream-kafka-consumer-filter/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-filter/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-filter/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java b/labx-11/labx-11-sc-stream-kafka-consumer-filter/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-filter/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-filter/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/resources/application.yml new file mode 100644 index 00000000..18a2a4d9 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-filter/src/main/resources/application.yml @@ -0,0 +1,22 @@ +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 # 消费者分组 + # 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/labx-11-sc-stream-kafka-consumer-partitioning/pom.xml b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/pom.xml new file mode 100644 index 00000000..438af790 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/pom.xml @@ -0,0 +1,58 @@ + + + + labx-11 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-11-sc-stream-kafka-consumer-partitioning + + + 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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/ConsumerApplication.java b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/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-partitioning/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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java new file mode 100644 index 00000000..1be848c7 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/Demo01Consumer.java @@ -0,0 +1,19 @@ +package cn.iocoder.springcloud.labx11.kafkademo.consumerdemo.listener; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.messaging.Message; +import org.springframework.stereotype.Component; + +@Component +public class Demo01Consumer { + + private Logger logger = LoggerFactory.getLogger(getClass()); + + @StreamListener(MySink.DEMO01_INPUT) + public void onMessage(Message message) { + logger.info("[onMessage][线程编号:{} 消息内容:{}]", Thread.currentThread().getId(), message); + } + +} diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/listener/MySink.java b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/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-partitioning/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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/consumerdemo/message/Demo01Message.java b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/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-partitioning/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-partitioning/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/src/main/resources/application.yml new file mode 100644 index 00000000..80cc70e9 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-consumer-partitioning/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/labx-11-sc-stream-kafka-producer-demo/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java b/labx-11/labx-11-sc-stream-kafka-producer-demo/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java index d159eeac..6f066581 100644 --- a/labx-11/labx-11-sc-stream-kafka-producer-demo/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java +++ b/labx-11/labx-11-sc-stream-kafka-producer-demo/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java @@ -34,4 +34,20 @@ public class Demo01Controller { return mySource.demo01Output().send(springMessage); } + @GetMapping("/send_tag") + public boolean sendTag() { + for (String tag : new String[]{"yunai", "yutou", "tudou"}) { + // 创建 Message + Demo01Message message = new Demo01Message() + .setId(new Random().nextInt()); + // 创建 Spring Message 对象 + Message springMessage = MessageBuilder.withPayload(message) + .setHeader("tag", tag) // 设置 Tag + .build(); + // 发送消息 + mySource.demo01Output().send(springMessage); + } + return true; + } + } diff --git a/labx-11/labx-11-sc-stream-kafka-producer-partitioning/pom.xml b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/pom.xml new file mode 100644 index 00000000..877a2cb7 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/pom.xml @@ -0,0 +1,58 @@ + + + + labx-11 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-11-sc-stream-kafka-producer-partitioning + + + 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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/ProducerApplication.java b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/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-partitioning/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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java new file mode 100644 index 00000000..900fdfe3 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java @@ -0,0 +1,41 @@ +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_orderly") + public boolean sendOrderly() { + // 发送 3 条相同 id 的消息 + int id = new Random().nextInt(); + for (int i = 0; i < 3; i++) { + // 创建 Message + Demo01Message message = new Demo01Message().setId(id); + // 创建 Spring Message 对象 + Message springMessage = MessageBuilder.withPayload(message) + .build(); + // 发送消息 + mySource.demo01Output().send(springMessage); + } + return true; + } + +} diff --git a/labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/Demo01Message.java b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/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-partitioning/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-partitioning/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/message/MySource.java b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/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-partitioning/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-partitioning/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/resources/application.yml new file mode 100644 index 00000000..0300d455 --- /dev/null +++ b/labx-11/labx-11-sc-stream-kafka-producer-partitioning/src/main/resources/application.yml @@ -0,0 +1,30 @@ +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 + # Producer 配置项,对应 ProducerProperties 类 + producer: + partition-key-expression: payload['id'] # 分区 key 表达式。该表达式基于 Spring EL,从消息中获得分区 key。 + # 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: + sync: true # 是否同步发送消息,默认为 false 异步。 + +server: + port: 18080 diff --git a/labx-11/pom.xml b/labx-11/pom.xml index fd418177..2d500581 100644 --- a/labx-11/pom.xml +++ b/labx-11/pom.xml @@ -26,6 +26,13 @@ labx-11-sc-stream-kafka-consumer-concurrency + + labx-11-sc-stream-kafka-producer-partitioning + labx-11-sc-stream-kafka-consumer-concurrency + labx-11-sc-stream-kafka-consumer-partitioning + + + labx-11-sc-stream-kafka-consumer-filter