diff --git a/spring-pulsar-sample-apps/src/main/java/app1/PulsarBootApp.java b/spring-pulsar-sample-apps/src/main/java/app1/PulsarBootApp.java deleted file mode 100644 index 2249968e..00000000 --- a/spring-pulsar-sample-apps/src/main/java/app1/PulsarBootApp.java +++ /dev/null @@ -1,85 +0,0 @@ -/* - * Copyright 2022 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package app1; - -import org.apache.pulsar.common.schema.SchemaType; - -import org.springframework.boot.ApplicationRunner; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.annotation.Bean; -import org.springframework.pulsar.annotation.PulsarListener; -import org.springframework.pulsar.core.PulsarTemplate; - -@SpringBootApplication -public class PulsarBootApp { - - public static void main(String[] args) { - SpringApplication.run(PulsarBootApp.class, args); - } - - @Bean - public ApplicationRunner runner(PulsarTemplate pulsarTemplate) { - String topic = "hello-pulsar-exclusive-2"; - return args -> { -// for (int i = 0; i < 100; i ++) { -// pulsarTemplate.send(topic, "This is message " + (i + 1)); -// } - Foo foo = new Foo(); - foo.setFoo("Foo"); - foo.setBar("Bar"); - pulsarTemplate.send(topic, foo); - - }; - } - - @PulsarListener(subscriptionName = "test-exclusive-sub-2", topics = "hello-pulsar-exclusive-2", schemaType = SchemaType.JSON) - public void listen(Foo foo) { - //... - } - - static class Foo { - String foo; - String bar; - - public String getFoo() { - return this.foo; - } - - public void setFoo(String foo) { - this.foo = foo; - } - - public String getBar() { - return this.bar; - } - - public void setBar(String bar) { - this.bar = bar; - } - - @Override - public String toString() { - return "Foo{" + - "foo='" + this.foo + '\'' + - ", bar='" + this.bar + '\'' + - '}'; - } - } - - -} diff --git a/spring-pulsar-sample-apps/src/main/java/app1/SpringPulsarBootApp.java b/spring-pulsar-sample-apps/src/main/java/app1/SpringPulsarBootApp.java new file mode 100644 index 00000000..a6ea82a2 --- /dev/null +++ b/spring-pulsar-sample-apps/src/main/java/app1/SpringPulsarBootApp.java @@ -0,0 +1,101 @@ +/* + * Copyright 2022 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package app1; + +import java.util.UUID; + +import org.apache.pulsar.common.schema.SchemaType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.springframework.boot.ApplicationRunner; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; +import org.springframework.pulsar.annotation.PulsarListener; +import org.springframework.pulsar.core.PulsarTemplate; + +@SpringBootApplication +public class SpringPulsarBootApp { + + Logger logger = LoggerFactory.getLogger(SpringPulsarBootApp.class); + + public static void main(String[] args) { + SpringApplication.run(SpringPulsarBootApp.class, args); + } + + @Bean + public ApplicationRunner runner1(PulsarTemplate pulsarTemplate) { + + String topic1 = "hello-pulsar-exclusive-1"; + + return args -> { + for (int i = 0; i < 10; i++) { + pulsarTemplate.send(topic1, "This is message " + (i + 1)); + } + }; + } + + @PulsarListener(subscriptionName = "subscription-1", topics = "hello-pulsar-exclusive-1") + public void listen1(String message) { + this.logger.info(message); + } + + @Bean + public ApplicationRunner runner2(PulsarTemplate pulsarTemplate) { + + String topic1 = "hello-pulsar-exclusive-2"; + + return args -> { + for (int i = 0; i < 10; i++) { + pulsarTemplate.send(topic1, i); + } + }; + } + + @PulsarListener(subscriptionName = "subscription-2", topics = "hello-pulsar-exclusive-2") + public void listen2(Integer message) { + this.logger.info("Message received :" + message); + } + + @Bean + public ApplicationRunner runner3(PulsarTemplate pulsarTemplate) { + + String topic = "hello-pulsar-exclusive-3"; + return args -> { + for (int i = 0; i < 10; i++) { + Foo foo = new Foo(i + "-" + "Foo-" + UUID.randomUUID(), i + "-" + "Bar-" + UUID.randomUUID()); + pulsarTemplate.send(topic, foo); + } + }; + } + + @PulsarListener(subscriptionName = "subscription-3", topics = "hello-pulsar-exclusive-3", schemaType = SchemaType.JSON) + public void listen3(Foo message) { + this.logger.info("Message received :" + message); + } + + record Foo(String foo, String bar) { + @Override + public String toString() { + return "Foo{" + + "foo='" + this.foo + '\'' + + ", bar='" + this.bar + '\'' + + '}'; + } + } +} diff --git a/spring-pulsar-sample-apps/src/main/java/app4/FailoverConsumerApp.java b/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java similarity index 87% rename from spring-pulsar-sample-apps/src/main/java/app4/FailoverConsumerApp.java rename to spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java index 8ae725be..b9d62f75 100644 --- a/spring-pulsar-sample-apps/src/main/java/app4/FailoverConsumerApp.java +++ b/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java @@ -14,13 +14,15 @@ * limitations under the License. */ -package app4; +package app2; import java.io.Serial; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.TopicMetadata; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.boot.ApplicationRunner; import org.springframework.boot.SpringApplication; @@ -32,11 +34,10 @@ import org.springframework.pulsar.core.PulsarTemplate; @SpringBootApplication public class FailoverConsumerApp { + Logger logger = LoggerFactory.getLogger(FailoverConsumerApp.class); + public static void main(String[] args) { - String[] args1 = new String[]{ -// "--spring.pulsar.consumer.subscription-type=Failover", - "--spring.pulsar.producer.messageRoutingMode=CustomPartition"}; - SpringApplication.run(FailoverConsumerApp.class, args1); + SpringApplication.run(FailoverConsumerApp.class, "--spring.pulsar.producer.messageRoutingMode=CustomPartition"); } @Bean @@ -49,23 +50,22 @@ public class FailoverConsumerApp { pulsarTemplate.sendAsync(topic, "hello buzz doe 2", new BuzzRouter()); Thread.sleep(1_000); } - System.exit(0); }; } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", subscriptionType = "failover") public void listen1(String foo) { - //... + this.logger.info("failover-listen1 : " + foo); } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", subscriptionType = "failover") public void listen2(String foo) { - //... + this.logger.info("failover-listen2 : " + foo); } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", subscriptionType = "failover") public void listen(String foo) { - //... + this.logger.info("failover-listen3 : " + foo); } static class FooRouter implements MessageRouter { @@ -97,5 +97,4 @@ public class FailoverConsumerApp { return 2; } } - } diff --git a/spring-pulsar-sample-apps/src/main/java/app2/ProducerApp.java b/spring-pulsar-sample-apps/src/main/java/app2/ProducerApp.java deleted file mode 100644 index 0750bec3..00000000 --- a/spring-pulsar-sample-apps/src/main/java/app2/ProducerApp.java +++ /dev/null @@ -1,88 +0,0 @@ -/* - * Copyright 2022 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package app2; - -import java.io.Serial; -import java.util.Random; - -import org.apache.pulsar.client.api.Message; -import org.apache.pulsar.client.api.MessageRouter; -import org.apache.pulsar.client.api.TopicMetadata; - -import org.springframework.boot.ApplicationRunner; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.annotation.Bean; -import org.springframework.pulsar.core.PulsarTemplate; - -@SpringBootApplication -public class ProducerApp { - - - public static void main(String[] args) { - String[] args1 = new String[]{ -// "--spring.pulsar.consumer.subscription-type=Failover", - "--spring.pulsar.producer.messageRoutingMode=CustomPartition"}; - SpringApplication.run(ProducerApp.class, args1); - } - - @Bean - public ApplicationRunner runner(PulsarTemplate pulsarTemplate) { - String topic = "failover-demo-topic"; - return args -> { - for (int i = 0; i < 100; i++) { - pulsarTemplate.sendAsync(topic, "hello john doex " + new Random().nextInt(), new FooRouter()); - pulsarTemplate.sendAsync(topic, "hello alice doex " + new Random().nextInt(), new BarRouter()); - if (i % 2 == 0) { - pulsarTemplate.sendAsync(topic, "hello buzz doex " + new Random().nextInt(), new BuzzRouter()); - } - Thread.sleep(5_000); - } - }; - } - - static class FooRouter implements MessageRouter { - @Serial - private static final long serialVersionUID = -1L; - - @Override - public int choosePartition(Message msg, TopicMetadata metadata) { - return 0; - } - } - - static class BarRouter implements MessageRouter { - @Serial - private static final long serialVersionUID = -1L; - - @Override - public int choosePartition(Message msg, TopicMetadata metadata) { - return 1; - } - } - - static class BuzzRouter implements MessageRouter { - @Serial - private static final long serialVersionUID = -1L; - - @Override - public int choosePartition(Message msg, TopicMetadata metadata) { - return 2; - } - } - -} diff --git a/spring-pulsar-sample-apps/src/main/java/app5/FailoverConsumer.java b/spring-pulsar-sample-apps/src/main/java/app5/FailoverConsumer.java deleted file mode 100644 index d5c5d976..00000000 --- a/spring-pulsar-sample-apps/src/main/java/app5/FailoverConsumer.java +++ /dev/null @@ -1,42 +0,0 @@ -/* - * Copyright 2022 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package app5; - -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.pulsar.annotation.PulsarListener; - -@SpringBootApplication -public class FailoverConsumer { - - - public static void main(String[] args) { - String[] args1 = new String[]{ -// "--spring.pulsar.consumer.subscription-type=Failover", - "--spring.pulsar.producer.messageRoutingMode=CustomPartition"}; - SpringApplication.run(FailoverConsumer.class, args1); - } - - - @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", subscriptionType = "shared") - public void listen1(String foo) { - //... - } - - -} - diff --git a/spring-pulsar-sample-apps/src/main/resources/application.yml b/spring-pulsar-sample-apps/src/main/resources/application.yml new file mode 100644 index 00000000..45c3b356 --- /dev/null +++ b/spring-pulsar-sample-apps/src/main/resources/application.yml @@ -0,0 +1,3 @@ +logging: + level: + org.apache.pulsar: error diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index 99ba7ed8..0675730e 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -39,7 +39,6 @@ import org.apache.pulsar.client.api.SubscriptionType; import org.springframework.context.ApplicationEventPublisher; import org.springframework.core.task.AsyncListenableTaskExecutor; import org.springframework.core.task.SimpleAsyncTaskExecutor; -import org.springframework.lang.Nullable; import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.event.ConsumerFailedToStartEvent; import org.springframework.pulsar.event.ConsumerStartedEvent; @@ -59,8 +58,6 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess private volatile boolean running = false; - private String beanName; - private volatile ListenableFuture listenerConsumerFuture; private volatile Listener listenerConsumer; @@ -135,16 +132,6 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess this.running = running; } - /** - * Return the bean name. - * - * @return the bean name. - */ - @Nullable - public String getBeanName() { - return this.beanName; - } - @Override public void destroy() {