Merge pull request #47707 from onobc

* pr/47707:
  Polish "Remove Spring Pulsar Reactive support"
  Remove Spring Pulsar Reactive support

Closes gh-47707
This commit is contained in:
Stéphane Nicoll
2025-10-20 18:15:45 +02:00
32 changed files with 542 additions and 2195 deletions
@@ -178,7 +178,6 @@ public class AntoraAsciidocAttributes {
addSpringDataDependencyVersion(attributes, internal, "spring-data-redis");
addSpringDataDependencyVersion(attributes, internal, "spring-data-rest", "spring-data-rest-core");
addSpringDataDependencyVersion(attributes, internal, "spring-data-ldap");
addDependencyVersion(attributes, "pulsar-client-reactive-api", "org.apache.pulsar:pulsar-client-reactive-api");
addDependencyVersion(attributes, "pulsar-client-api", "org.apache.pulsar:pulsar-client-api");
}
@@ -36,7 +36,6 @@ url-native-build-tools-docs-maven-plugin={url-native-build-tools-docs}/maven-plu
url-paketo-docs=https://paketo.io/docs
url-paketo-docs-java-buildpack={url-paketo-docs}/buildpacks/language-family-buildpacks/java
url-pulsar-client-api-javadoc=https://javadoc.io/doc/org.apache.pulsar/pulsar-client-api/{version-pulsar-client-api}
url-pulsar-client-reactive-api-javadoc=https://javadoc.io/doc/org.apache.pulsar/pulsar-client-reactive-api/{version-pulsar-client-reactive-api}
url-spring-boot-for-apache-geode-docs=https://docs.spring.io/spring-boot-data-geode-build/2.0.x/reference/html5
url-spring-boot-for-apache-geode-site=https://github.com/spring-projects/spring-boot-data-geode
url-spring-data-cassandra-docs=https://docs.spring.io/spring-data/cassandra/reference/{antoraversion-spring-data-cassandra}
@@ -87,7 +86,6 @@ url-jackson2-databind-javadoc=https://javadoc.io/doc/com.fasterxml.jackson.core/
javadoc-location-com-fasterxml-jackson-annotation={url-jackson-annotations-javadoc}
javadoc-location-com-fasterxml-jackson-databind={url-jackson2-databind-javadoc}
javadoc-location-org-apache-pulsar-client-api={url-pulsar-client-api-javadoc}
javadoc-location-org-apache-pulsar-reactive-client-api={url-pulsar-client-reactive-api-javadoc}
javadoc-location-org-springframework-data-cassandra={url-spring-data-cassandra-javadoc}
javadoc-location-org-springframework-data-convert={url-spring-data-commons-javadoc}
javadoc-location-org-springframework-data-querydsl={url-spring-data-commons-javadoc}
@@ -132,4 +130,3 @@ code-spring-boot-latest=https://github.com/{github-repo}/tree/main
code-spring-boot-servlet-src={code-spring-boot}/module/spring-boot-servlet/src/main/java/org/springframework/boot/servlet
code-spring-boot-thymeleaf-src={code-spring-boot}/module/spring-boot-thymeleaf/src/main/java/org/springframework/boot/thymeleaf
code-spring-boot-webmvc-src={code-spring-boot}/module/spring-boot-webmvc/src/main/java/org/springframework/boot/webmvc
@@ -287,7 +287,6 @@ class AntoraAsciidocAttributesTests {
addMockJacksonCoreVersion(versions, "jackson-databind", version);
addMockJacksonCoreVersion(versions, "jackson-databind", version);
versions.put("org.apache.pulsar:pulsar-client-api", version);
versions.put("org.apache.pulsar:pulsar-client-reactive-api", version);
versions.put("tools.jackson.dataformat:jackson-dataformat-xml", version);
return versions;
}
@@ -204,7 +204,6 @@ dependencies {
implementation("org.springframework.kafka:spring-kafka")
implementation("org.springframework.kafka:spring-kafka-test")
implementation("org.springframework.pulsar:spring-pulsar")
implementation("org.springframework.pulsar:spring-pulsar-reactive")
implementation("org.springframework.restdocs:spring-restdocs-mockmvc")
implementation("org.springframework.restdocs:spring-restdocs-webtestclient")
implementation("org.springframework.security:spring-security-config")
@@ -1781,15 +1781,11 @@
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.additional-properties[#messaging.pulsar.additional-properties]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.admin.auth[#messaging.pulsar.admin.auth]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.admin[#messaging.pulsar.admin]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.connecting-reactive[#messaging.pulsar.connecting-reactive]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.connecting.auth[#messaging.pulsar.connecting.auth]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.connecting.ssl[#messaging.pulsar.connecting.ssl]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.connecting[#messaging.pulsar.connecting]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.reading-reactive[#messaging.pulsar.reading-reactive]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.reading[#messaging.pulsar.reading]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.receiving-reactive[#messaging.pulsar.receiving-reactive]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.receiving[#messaging.pulsar.receiving]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.sending-reactive[#messaging.pulsar.sending-reactive]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar.sending[#messaging.pulsar.sending]
* xref:reference:messaging/pulsar.adoc#messaging.pulsar[#messaging.pulsar]
* xref:reference:messaging/rsocket.adoc#messaging.rsocket.messaging[#boot-features-rsocket-messaging]
@@ -3,10 +3,9 @@
https://pulsar.apache.org/[Apache Pulsar] is supported by providing auto-configuration of the {url-spring-pulsar-site}[Spring for Apache Pulsar] project.
Spring Boot will auto-configure and register the classic (imperative) Spring for Apache Pulsar components when `org.springframework.pulsar:spring-pulsar` is on the classpath.
It will do the same for the reactive components when `org.springframework.pulsar:spring-pulsar-reactive` is on the classpath.
Spring Boot will auto-configure and register the Spring for Apache Pulsar components when `org.springframework.pulsar:spring-pulsar` is on the classpath.
There are `spring-boot-starter-pulsar` and `spring-boot-starter-pulsar-reactive` starters for conveniently collecting the dependencies for imperative and reactive use, respectively.
There is the `spring-boot-starter-pulsar` starter for conveniently collecting the dependencies for use.
@@ -68,14 +67,6 @@ You can follow {url-spring-pulsar-docs}/reference/pulsar/pulsar-client.html#tls-
For complete details on the client and authentication see the Spring for Apache Pulsar {url-spring-pulsar-docs}/reference/pulsar/pulsar-client.html[reference documentation].
[[messaging.pulsar.connecting-reactive]]
== Connecting to Pulsar Reactively
When the Reactive auto-configuration is activated, Spring Boot will auto-configure and register a javadoc:org.apache.pulsar.reactive.client.api.ReactivePulsarClient[] bean.
The javadoc:org.apache.pulsar.reactive.client.api.ReactivePulsarClient[] adapts an instance of the previously described javadoc:org.apache.pulsar.client.api.PulsarClient[].
Therefore, follow the previous section to configure the javadoc:org.apache.pulsar.client.api.PulsarClient[] used by the javadoc:org.apache.pulsar.reactive.client.api.ReactivePulsarClient[].
[[messaging.pulsar.admin]]
@@ -120,25 +111,6 @@ If you need more control over the message being sent, you can pass in a javadoc:
[[messaging.pulsar.sending-reactive]]
== Sending a Message Reactively
When the Reactive auto-configuration is activated, Spring's javadoc:org.springframework.pulsar.reactive.core.ReactivePulsarTemplate[] is auto-configured, and you can use it to send messages, as shown in the following example:
include-code::MyBean[]
The javadoc:org.springframework.pulsar.reactive.core.ReactivePulsarTemplate[] relies on a javadoc:org.springframework.pulsar.reactive.core.ReactivePulsarSenderFactory[] to actually create the underlying sender.
Spring Boot auto-configuration also provides this sender factory, which by default, caches the producers that it creates.
You can configure the sender factory and cache settings by specifying any of the `spring.pulsar.producer.\*` and `spring.pulsar.producer.cache.*` prefixed application properties.
If you need more control over the sender factory configuration, consider registering one or more javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageSenderBuilderCustomizer[] beans.
These customizers are applied to all created senders.
You can also pass in a javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageSenderBuilderCustomizer[] when sending a message to only affect the current sender.
If you need more control over the message being sent, you can pass in a javadoc:org.springframework.pulsar.reactive.core.MessageSpecBuilderCustomizer[] when sending a message.
[[messaging.pulsar.receiving]]
== Receiving a Message
@@ -156,22 +128,7 @@ You can also customize a single listener by setting the `consumerCustomizer` att
If you need more control over the actual container factory configuration, consider registering one or more `PulsarContainerFactoryCustomizer<ConcurrentPulsarListenerContainerFactory<?>>` beans.
[[messaging.pulsar.receiving-reactive]]
== Receiving a Message Reactively
When the Apache Pulsar infrastructure is present and the Reactive auto-configuration is activated, any bean can be annotated with javadoc:org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener[format=annotation] to create a reactive listener endpoint.
The following component creates a reactive listener endpoint on the `someTopic` topic:
include-code::MyBean[]
Spring Boot auto-configuration provides all the components necessary for javadoc:org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener[], such as the javadoc:org.springframework.pulsar.reactive.config.ReactivePulsarListenerContainerFactory[] and the consumer factory it uses to construct the underlying reactive Pulsar consumers.
You can configure these components by specifying any of the `spring.pulsar.listener.\*` and `spring.pulsar.consumer.*` prefixed application properties.
If you need more control over the configuration of the consumer factory, consider registering one or more javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageConsumerBuilderCustomizer[] beans.
These customizers are applied to all consumers created by the factory, and therefore all javadoc:org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener[format=annotation] instances.
You can also customize a single listener by setting the `consumerCustomizer` attribute of the javadoc:org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener[format=annotation] annotation.
If you need more control over the actual container factory configuration, consider registering one or more `PulsarContainerFactoryCustomizer<DefaultReactivePulsarListenerContainerFactory<?>>` beans.
[[messaging.pulsar.reading]]
== Reading a Message
@@ -193,23 +150,6 @@ You can also customize a single listener by setting the `readerCustomizer` attri
If you need more control over the actual container factory configuration, consider registering one or more `PulsarContainerFactoryCustomizer<DefaultPulsarReaderContainerFactory<?>>` beans.
[[messaging.pulsar.reading-reactive]]
== Reading a Message Reactively
When the Apache Pulsar infrastructure is present and the Reactive auto-configuration is activated, Spring's javadoc:org.springframework.pulsar.reactive.core.ReactivePulsarReaderFactory[] is provided, and you can use it to create a reader in order to read messages in a reactive fashion.
The following component creates a reader using the provided factory and reads a single message from 5 minutes ago from the `someTopic` topic:
include-code::MyBean[]
Spring Boot auto-configuration provides this reader factory which can be customized by setting any of the `spring.pulsar.reader.*` prefixed application properties.
If you need more control over the reader factory configuration, consider passing in one or more javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer[] instances when using the factory to create a reader.
If you need more control over the reader factory configuration, consider registering one or more javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer[] beans.
These customizers are applied to all created readers.
You can also pass one or more javadoc:org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer[] when creating a reader to only apply the customizations to the created reader.
TIP: For more details on any of the above components and to discover other available features, see the Spring for Apache Pulsar {url-spring-pulsar-docs}[reference documentation].
@@ -219,8 +159,6 @@ TIP: For more details on any of the above components and to discover other avail
Spring for Apache Pulsar supports transactions when using javadoc:org.springframework.pulsar.core.PulsarTemplate[] and javadoc:org.springframework.pulsar.annotation.PulsarListener[format=annotation].
NOTE: Transactions are not currently supported when using the reactive variants.
Setting the configprop:spring.pulsar.transaction.enabled[] property to `true` will:
* Configure a javadoc:org.springframework.pulsar.transaction.PulsarTransactionManager[] bean
@@ -1,51 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.readingreactive;
import java.time.Instant;
import java.util.List;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.reactive.client.api.StartAtSpec;
import reactor.core.publisher.Mono;
import org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactivePulsarReaderFactory;
import org.springframework.stereotype.Component;
@Component
public class MyBean {
private final ReactivePulsarReaderFactory<String> pulsarReaderFactory;
public MyBean(ReactivePulsarReaderFactory<String> pulsarReaderFactory) {
this.pulsarReaderFactory = pulsarReaderFactory;
}
@SuppressWarnings("unused")
public void someMethod() {
ReactiveMessageReaderBuilderCustomizer<String> readerBuilderCustomizer = (readerBuilder) -> readerBuilder
.topic("someTopic")
.startAtSpec(StartAtSpec.ofInstant(Instant.now().minusSeconds(5)));
Mono<Message<String>> message = this.pulsarReaderFactory
.createReader(Schema.STRING, List.of(readerBuilderCustomizer))
.readOne();
// ...
}
}
@@ -1,33 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.receivingreactive;
import reactor.core.publisher.Mono;
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener;
import org.springframework.stereotype.Component;
@Component
public class MyBean {
@ReactivePulsarListener(topics = "someTopic")
public Mono<Void> processMessage(String content) {
// ...
return Mono.empty();
}
}
@@ -1,35 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.sendingreactive;
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate;
import org.springframework.stereotype.Component;
@Component
public class MyBean {
private final ReactivePulsarTemplate<String> pulsarTemplate;
public MyBean(ReactivePulsarTemplate<String> pulsarTemplate) {
this.pulsarTemplate = pulsarTemplate;
}
public void someMethod() {
this.pulsarTemplate.send("someTopic", "Hello").subscribe();
}
}
@@ -1,45 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.readingreactive
import org.apache.pulsar.client.api.Schema
import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderBuilder
import org.apache.pulsar.reactive.client.api.StartAtSpec
import org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer
import org.springframework.pulsar.reactive.core.ReactivePulsarReaderFactory
import org.springframework.stereotype.Component
import java.time.Instant
@Suppress("UNUSED_PARAMETER", "UNUSED_VARIABLE")
@Component
class MyBean(private val pulsarReaderFactory: ReactivePulsarReaderFactory<String>) {
fun someMethod() {
val readerBuilderCustomizer = ReactiveMessageReaderBuilderCustomizer {
readerBuilder: ReactiveMessageReaderBuilder<String> ->
readerBuilder
.topic("someTopic")
.startAtSpec(StartAtSpec.ofInstant(Instant.now().minusSeconds(5)))
}
val message = pulsarReaderFactory
.createReader(Schema.STRING, listOf(readerBuilderCustomizer))
.readOne()
// ...
}
}
@@ -1,33 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.receivingreactive
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener
import org.springframework.stereotype.Component
import reactor.core.publisher.Mono
@Component
@Suppress("UNUSED_PARAMETER")
class MyBean {
@ReactivePulsarListener(topics = ["someTopic"])
fun processMessage(content: String?): Mono<Void> {
// ...
return Mono.empty()
}
}
@@ -1,29 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.docs.messaging.pulsar.sendingreactive
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate
import org.springframework.stereotype.Component
@Component
class MyBean(private val pulsarTemplate: ReactivePulsarTemplate<String>) {
fun someMethod() {
pulsarTemplate.send("someTopic", "Hello").subscribe()
}
}
-1
View File
@@ -32,7 +32,6 @@ dependencies {
optional(project(":core:spring-boot-autoconfigure"))
optional(project(":core:spring-boot-docker-compose"))
optional(project(":core:spring-boot-testcontainers"))
optional("org.springframework.pulsar:spring-pulsar-reactive")
optional("org.testcontainers:testcontainers-pulsar")
dockerTestImplementation(project(":core:spring-boot-test"))
@@ -27,7 +27,6 @@ import org.testcontainers.pulsar.PulsarContainer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.boot.pulsar.autoconfigure.PulsarAutoConfiguration;
import org.springframework.boot.pulsar.autoconfigure.PulsarReactiveAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.testsupport.container.TestImage;
import org.springframework.context.annotation.Configuration;
@@ -76,7 +75,7 @@ class PulsarAutoConfigurationIntegrationTests {
}
@Configuration(proxyBeanMethods = false)
@ImportAutoConfiguration({ PulsarAutoConfiguration.class, PulsarReactiveAutoConfiguration.class })
@ImportAutoConfiguration(PulsarAutoConfiguration.class)
@Import(TestService.class)
static class TestConfiguration {
@@ -19,23 +19,33 @@ package org.springframework.boot.pulsar.autoconfigure;
import java.util.ArrayList;
import java.util.List;
import org.apache.pulsar.client.admin.PulsarAdminBuilder;
import org.apache.pulsar.client.api.ClientBuilder;
import org.apache.pulsar.client.api.ConsumerBuilder;
import org.apache.pulsar.client.api.ProducerBuilder;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.ReaderBuilder;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.interceptor.ProducerInterceptor;
import org.apache.pulsar.common.naming.TopicDomain;
import org.apache.pulsar.common.schema.SchemaType;
import org.jspecify.annotations.Nullable;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Defaults.SchemaInfo;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Defaults.TypeMapping;
import org.springframework.boot.thread.Threading;
import org.springframework.boot.util.LambdaSafe;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.context.annotation.Scope;
import org.springframework.core.env.Environment;
import org.springframework.core.task.VirtualThreadTaskExecutor;
import org.springframework.pulsar.annotation.EnablePulsar;
@@ -44,10 +54,17 @@ import org.springframework.pulsar.config.DefaultPulsarReaderContainerFactory;
import org.springframework.pulsar.config.PulsarAnnotationSupportBeanNames;
import org.springframework.pulsar.core.CachingPulsarProducerFactory;
import org.springframework.pulsar.core.ConsumerBuilderCustomizer;
import org.springframework.pulsar.core.DefaultPulsarClientFactory;
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
import org.springframework.pulsar.core.DefaultPulsarReaderFactory;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.DefaultTopicResolver;
import org.springframework.pulsar.core.ProducerBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdministration;
import org.springframework.pulsar.core.PulsarClientBuilderCustomizer;
import org.springframework.pulsar.core.PulsarClientFactory;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.core.PulsarProducerFactory;
import org.springframework.pulsar.core.PulsarReaderFactory;
@@ -55,11 +72,17 @@ import org.springframework.pulsar.core.PulsarTemplate;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.ReaderBuilderCustomizer;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.function.PulsarFunction;
import org.springframework.pulsar.function.PulsarFunctionAdministration;
import org.springframework.pulsar.function.PulsarSink;
import org.springframework.pulsar.function.PulsarSource;
import org.springframework.pulsar.listener.PulsarContainerProperties;
import org.springframework.pulsar.reader.PulsarReaderContainerProperties;
import org.springframework.pulsar.transaction.PulsarAwareTransactionManager;
import org.springframework.pulsar.transaction.PulsarTransactionManager;
import org.springframework.util.Assert;
/**
* {@link EnableAutoConfiguration Auto-configuration} for Apache Pulsar.
@@ -73,7 +96,7 @@ import org.springframework.pulsar.transaction.PulsarTransactionManager;
*/
@AutoConfiguration
@ConditionalOnClass({ PulsarClient.class, PulsarTemplate.class })
@Import(PulsarConfiguration.class)
@EnableConfigurationProperties(PulsarProperties.class)
public final class PulsarAutoConfiguration {
private final PulsarProperties properties;
@@ -85,6 +108,136 @@ public final class PulsarAutoConfiguration {
this.propertiesMapper = new PulsarPropertiesMapper(properties);
}
@Bean
@ConditionalOnMissingBean(PulsarConnectionDetails.class)
PropertiesPulsarConnectionDetails pulsarConnectionDetails() {
return new PropertiesPulsarConnectionDetails(this.properties);
}
@Bean
@ConditionalOnMissingBean(PulsarClientFactory.class)
DefaultPulsarClientFactory pulsarClientFactory(PulsarConnectionDetails connectionDetails,
ObjectProvider<PulsarClientBuilderCustomizer> customizersProvider) {
List<PulsarClientBuilderCustomizer> allCustomizers = new ArrayList<>();
allCustomizers.add((builder) -> this.propertiesMapper.customizeClientBuilder(builder, connectionDetails));
allCustomizers.addAll(customizersProvider.orderedStream().toList());
DefaultPulsarClientFactory clientFactory = new DefaultPulsarClientFactory(
(clientBuilder) -> applyClientBuilderCustomizers(allCustomizers, clientBuilder));
return clientFactory;
}
private void applyClientBuilderCustomizers(List<PulsarClientBuilderCustomizer> customizers,
ClientBuilder clientBuilder) {
customizers.forEach((customizer) -> customizer.customize(clientBuilder));
}
@Bean
@ConditionalOnMissingBean
PulsarClient pulsarClient(PulsarClientFactory clientFactory) {
return clientFactory.createClient();
}
@Bean
@ConditionalOnMissingBean
PulsarAdministration pulsarAdministration(PulsarConnectionDetails connectionDetails,
ObjectProvider<PulsarAdminBuilderCustomizer> pulsarAdminBuilderCustomizers) {
List<PulsarAdminBuilderCustomizer> allCustomizers = new ArrayList<>();
allCustomizers.add((builder) -> this.propertiesMapper.customizeAdminBuilder(builder, connectionDetails));
allCustomizers.addAll(pulsarAdminBuilderCustomizers.orderedStream().toList());
return new PulsarAdministration((adminBuilder) -> applyAdminBuilderCustomizers(allCustomizers, adminBuilder));
}
private void applyAdminBuilderCustomizers(List<PulsarAdminBuilderCustomizer> customizers,
PulsarAdminBuilder adminBuilder) {
customizers.forEach((customizer) -> customizer.customize(adminBuilder));
}
@Bean
@ConditionalOnMissingBean(SchemaResolver.class)
DefaultSchemaResolver pulsarSchemaResolver(ObjectProvider<SchemaResolverCustomizer<?>> schemaResolverCustomizers) {
DefaultSchemaResolver schemaResolver = new DefaultSchemaResolver();
addCustomSchemaMappings(schemaResolver, this.properties.getDefaults().getTypeMappings());
applySchemaResolverCustomizers(schemaResolverCustomizers.orderedStream().toList(), schemaResolver);
return schemaResolver;
}
private void addCustomSchemaMappings(DefaultSchemaResolver schemaResolver,
@Nullable List<TypeMapping> typeMappings) {
if (typeMappings != null) {
typeMappings.forEach((typeMapping) -> addCustomSchemaMapping(schemaResolver, typeMapping));
}
}
private void addCustomSchemaMapping(DefaultSchemaResolver schemaResolver, TypeMapping typeMapping) {
SchemaInfo schemaInfo = typeMapping.schemaInfo();
if (schemaInfo != null) {
Class<?> messageType = typeMapping.messageType();
SchemaType schemaType = schemaInfo.schemaType();
Class<?> messageKeyType = schemaInfo.messageKeyType();
Schema<?> schema = getSchema(schemaResolver, schemaType, messageType, messageKeyType);
schemaResolver.addCustomSchemaMapping(typeMapping.messageType(), schema);
}
}
private Schema<Object> getSchema(DefaultSchemaResolver schemaResolver, SchemaType schemaType, Class<?> messageType,
@Nullable Class<?> messageKeyType) {
Schema<Object> schema = schemaResolver.resolveSchema(schemaType, messageType, messageKeyType).orElseThrow();
Assert.state(schema != null, "'schema' must not be null");
return schema;
}
@SuppressWarnings("unchecked")
private void applySchemaResolverCustomizers(List<SchemaResolverCustomizer<?>> customizers,
DefaultSchemaResolver schemaResolver) {
LambdaSafe.callbacks(SchemaResolverCustomizer.class, customizers, schemaResolver)
.invoke((customizer) -> customizer.customize(schemaResolver));
}
@Bean
@ConditionalOnMissingBean(TopicResolver.class)
DefaultTopicResolver pulsarTopicResolver() {
DefaultTopicResolver topicResolver = new DefaultTopicResolver();
List<TypeMapping> typeMappings = this.properties.getDefaults().getTypeMappings();
if (typeMappings != null) {
typeMappings.forEach((typeMapping) -> addCustomTopicMapping(topicResolver, typeMapping));
}
return topicResolver;
}
private void addCustomTopicMapping(DefaultTopicResolver topicResolver, TypeMapping typeMapping) {
String topicName = typeMapping.topicName();
if (topicName != null) {
topicResolver.addCustomTopicMapping(typeMapping.messageType(), topicName);
}
}
@Bean
@ConditionalOnMissingBean
@ConditionalOnBooleanProperty(name = "spring.pulsar.function.enabled", matchIfMissing = true)
PulsarFunctionAdministration pulsarFunctionAdministration(PulsarAdministration pulsarAdministration,
ObjectProvider<PulsarFunction> pulsarFunctions, ObjectProvider<PulsarSink> pulsarSinks,
ObjectProvider<PulsarSource> pulsarSources) {
PulsarProperties.Function properties = this.properties.getFunction();
return new PulsarFunctionAdministration(pulsarAdministration, pulsarFunctions, pulsarSinks, pulsarSources,
properties.isFailFast(), properties.isPropagateFailures(), properties.isPropagateStopFailures());
}
@Bean
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@ConditionalOnMissingBean
@ConditionalOnBooleanProperty(name = "spring.pulsar.defaults.topic.enabled", matchIfMissing = true)
PulsarTopicBuilder pulsarTopicBuilder() {
return new PulsarTopicBuilder(TopicDomain.persistent, this.properties.getDefaults().getTopic().getTenant(),
this.properties.getDefaults().getTopic().getNamespace());
}
@Bean
@ConditionalOnMissingBean
PulsarContainerFactoryCustomizers pulsarContainerFactoryCustomizers(
ObjectProvider<PulsarContainerFactoryCustomizer<?>> customizers) {
return new PulsarContainerFactoryCustomizers(customizers.orderedStream().toList());
}
@Bean
@ConditionalOnMissingBean(PulsarProducerFactory.class)
@ConditionalOnBooleanProperty(name = "spring.pulsar.producer.cache.enabled", havingValue = false)
@@ -1,209 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.util.ArrayList;
import java.util.List;
import org.apache.pulsar.client.admin.PulsarAdminBuilder;
import org.apache.pulsar.client.api.ClientBuilder;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.common.naming.TopicDomain;
import org.apache.pulsar.common.schema.SchemaType;
import org.jspecify.annotations.Nullable;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Defaults.SchemaInfo;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Defaults.TypeMapping;
import org.springframework.boot.util.LambdaSafe;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Scope;
import org.springframework.pulsar.core.DefaultPulsarClientFactory;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.DefaultTopicResolver;
import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdministration;
import org.springframework.pulsar.core.PulsarClientBuilderCustomizer;
import org.springframework.pulsar.core.PulsarClientFactory;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.function.PulsarFunction;
import org.springframework.pulsar.function.PulsarFunctionAdministration;
import org.springframework.pulsar.function.PulsarSink;
import org.springframework.pulsar.function.PulsarSource;
import org.springframework.util.Assert;
/**
* Common configuration used by both {@link PulsarAutoConfiguration} and
* {@link PulsarReactiveAutoConfiguration}. A separate configuration class is used so that
* {@link PulsarAutoConfiguration} can be excluded for reactive only application.
*
* @author Chris Bono
* @author Phillip Webb
*/
@Configuration(proxyBeanMethods = false)
@EnableConfigurationProperties(PulsarProperties.class)
class PulsarConfiguration {
private final PulsarProperties properties;
private final PulsarPropertiesMapper propertiesMapper;
PulsarConfiguration(PulsarProperties properties) {
this.properties = properties;
this.propertiesMapper = new PulsarPropertiesMapper(properties);
}
@Bean
@ConditionalOnMissingBean(PulsarConnectionDetails.class)
PropertiesPulsarConnectionDetails pulsarConnectionDetails() {
return new PropertiesPulsarConnectionDetails(this.properties);
}
@Bean
@ConditionalOnMissingBean(PulsarClientFactory.class)
DefaultPulsarClientFactory pulsarClientFactory(PulsarConnectionDetails connectionDetails,
ObjectProvider<PulsarClientBuilderCustomizer> customizersProvider) {
List<PulsarClientBuilderCustomizer> allCustomizers = new ArrayList<>();
allCustomizers.add((builder) -> this.propertiesMapper.customizeClientBuilder(builder, connectionDetails));
allCustomizers.addAll(customizersProvider.orderedStream().toList());
DefaultPulsarClientFactory clientFactory = new DefaultPulsarClientFactory(
(clientBuilder) -> applyClientBuilderCustomizers(allCustomizers, clientBuilder));
return clientFactory;
}
private void applyClientBuilderCustomizers(List<PulsarClientBuilderCustomizer> customizers,
ClientBuilder clientBuilder) {
customizers.forEach((customizer) -> customizer.customize(clientBuilder));
}
@Bean
@ConditionalOnMissingBean
PulsarClient pulsarClient(PulsarClientFactory clientFactory) {
return clientFactory.createClient();
}
@Bean
@ConditionalOnMissingBean
PulsarAdministration pulsarAdministration(PulsarConnectionDetails connectionDetails,
ObjectProvider<PulsarAdminBuilderCustomizer> pulsarAdminBuilderCustomizers) {
List<PulsarAdminBuilderCustomizer> allCustomizers = new ArrayList<>();
allCustomizers.add((builder) -> this.propertiesMapper.customizeAdminBuilder(builder, connectionDetails));
allCustomizers.addAll(pulsarAdminBuilderCustomizers.orderedStream().toList());
return new PulsarAdministration((adminBuilder) -> applyAdminBuilderCustomizers(allCustomizers, adminBuilder));
}
private void applyAdminBuilderCustomizers(List<PulsarAdminBuilderCustomizer> customizers,
PulsarAdminBuilder adminBuilder) {
customizers.forEach((customizer) -> customizer.customize(adminBuilder));
}
@Bean
@ConditionalOnMissingBean(SchemaResolver.class)
DefaultSchemaResolver pulsarSchemaResolver(ObjectProvider<SchemaResolverCustomizer<?>> schemaResolverCustomizers) {
DefaultSchemaResolver schemaResolver = new DefaultSchemaResolver();
addCustomSchemaMappings(schemaResolver, this.properties.getDefaults().getTypeMappings());
applySchemaResolverCustomizers(schemaResolverCustomizers.orderedStream().toList(), schemaResolver);
return schemaResolver;
}
private void addCustomSchemaMappings(DefaultSchemaResolver schemaResolver,
@Nullable List<TypeMapping> typeMappings) {
if (typeMappings != null) {
typeMappings.forEach((typeMapping) -> addCustomSchemaMapping(schemaResolver, typeMapping));
}
}
private void addCustomSchemaMapping(DefaultSchemaResolver schemaResolver, TypeMapping typeMapping) {
SchemaInfo schemaInfo = typeMapping.schemaInfo();
if (schemaInfo != null) {
Class<?> messageType = typeMapping.messageType();
SchemaType schemaType = schemaInfo.schemaType();
Class<?> messageKeyType = schemaInfo.messageKeyType();
Schema<?> schema = getSchema(schemaResolver, schemaType, messageType, messageKeyType);
schemaResolver.addCustomSchemaMapping(typeMapping.messageType(), schema);
}
}
private Schema<Object> getSchema(DefaultSchemaResolver schemaResolver, SchemaType schemaType, Class<?> messageType,
@Nullable Class<?> messageKeyType) {
Schema<Object> schema = schemaResolver.resolveSchema(schemaType, messageType, messageKeyType).orElseThrow();
Assert.state(schema != null, "'schema' must not be null");
return schema;
}
@SuppressWarnings("unchecked")
private void applySchemaResolverCustomizers(List<SchemaResolverCustomizer<?>> customizers,
DefaultSchemaResolver schemaResolver) {
LambdaSafe.callbacks(SchemaResolverCustomizer.class, customizers, schemaResolver)
.invoke((customizer) -> customizer.customize(schemaResolver));
}
@Bean
@ConditionalOnMissingBean(TopicResolver.class)
DefaultTopicResolver pulsarTopicResolver() {
DefaultTopicResolver topicResolver = new DefaultTopicResolver();
List<TypeMapping> typeMappings = this.properties.getDefaults().getTypeMappings();
if (typeMappings != null) {
typeMappings.forEach((typeMapping) -> addCustomTopicMapping(topicResolver, typeMapping));
}
return topicResolver;
}
private void addCustomTopicMapping(DefaultTopicResolver topicResolver, TypeMapping typeMapping) {
String topicName = typeMapping.topicName();
if (topicName != null) {
topicResolver.addCustomTopicMapping(typeMapping.messageType(), topicName);
}
}
@Bean
@ConditionalOnMissingBean
@ConditionalOnBooleanProperty(name = "spring.pulsar.function.enabled", matchIfMissing = true)
PulsarFunctionAdministration pulsarFunctionAdministration(PulsarAdministration pulsarAdministration,
ObjectProvider<PulsarFunction> pulsarFunctions, ObjectProvider<PulsarSink> pulsarSinks,
ObjectProvider<PulsarSource> pulsarSources) {
PulsarProperties.Function properties = this.properties.getFunction();
return new PulsarFunctionAdministration(pulsarAdministration, pulsarFunctions, pulsarSinks, pulsarSources,
properties.isFailFast(), properties.isPropagateFailures(), properties.isPropagateStopFailures());
}
@Bean
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@ConditionalOnMissingBean
@ConditionalOnBooleanProperty(name = "spring.pulsar.defaults.topic.enabled", matchIfMissing = true)
PulsarTopicBuilder pulsarTopicBuilder() {
return new PulsarTopicBuilder(TopicDomain.persistent, this.properties.getDefaults().getTopic().getTenant(),
this.properties.getDefaults().getTopic().getNamespace());
}
@Bean
@ConditionalOnMissingBean
PulsarContainerFactoryCustomizers pulsarContainerFactoryCustomizers(
ObjectProvider<PulsarContainerFactoryCustomizer<?>> customizers) {
return new PulsarContainerFactoryCustomizers(customizers.orderedStream().toList());
}
}
@@ -1,218 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory;
import org.apache.pulsar.reactive.client.adapter.ProducerCacheProvider;
import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache;
import org.apache.pulsar.reactive.client.api.ReactivePulsarClient;
import org.apache.pulsar.reactive.client.producercache.CaffeineShadedProducerCacheProvider;
import org.jspecify.annotations.Nullable;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.util.LambdaSafe;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.pulsar.config.PulsarAnnotationSupportBeanNames;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory;
import org.springframework.pulsar.reactive.config.annotation.EnableReactivePulsar;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarConsumerFactory;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarReaderFactory;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarSenderFactory;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarSenderFactory.Builder;
import org.springframework.pulsar.reactive.core.ReactiveMessageConsumerBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactiveMessageSenderBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactivePulsarConsumerFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarReaderFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarSenderFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate;
import org.springframework.pulsar.reactive.listener.ReactivePulsarContainerProperties;
/**
* {@link EnableAutoConfiguration Auto-configuration} for Spring for Apache Pulsar
* Reactive.
*
* @author Chris Bono
* @author Christophe Bornet
* @since 4.0.0
*/
@AutoConfiguration(after = PulsarAutoConfiguration.class)
@ConditionalOnClass({ PulsarClient.class, ReactivePulsarClient.class, ReactivePulsarTemplate.class })
@Import(PulsarConfiguration.class)
public final class PulsarReactiveAutoConfiguration {
private final PulsarProperties properties;
private final PulsarReactivePropertiesMapper propertiesMapper;
PulsarReactiveAutoConfiguration(PulsarProperties properties) {
this.properties = properties;
this.propertiesMapper = new PulsarReactivePropertiesMapper(properties);
}
@Bean
@ConditionalOnMissingBean
ReactivePulsarClient reactivePulsarClient(PulsarClient pulsarClient) {
return AdaptedReactivePulsarClientFactory.create(pulsarClient);
}
@Bean
@ConditionalOnMissingBean(ProducerCacheProvider.class)
@ConditionalOnClass(CaffeineShadedProducerCacheProvider.class)
@ConditionalOnBooleanProperty(name = "spring.pulsar.producer.cache.enabled", matchIfMissing = true)
CaffeineShadedProducerCacheProvider reactivePulsarProducerCacheProvider() {
PulsarProperties.Producer.Cache properties = this.properties.getProducer().getCache();
return new CaffeineShadedProducerCacheProvider(properties.getExpireAfterAccess(), Duration.ofMinutes(10),
properties.getMaximumSize(), properties.getInitialCapacity());
}
@Bean
@ConditionalOnMissingBean
@ConditionalOnBooleanProperty(name = "spring.pulsar.producer.cache.enabled", matchIfMissing = true)
ReactiveMessageSenderCache reactivePulsarMessageSenderCache(
ObjectProvider<ProducerCacheProvider> producerCacheProvider) {
return reactivePulsarMessageSenderCache(producerCacheProvider.getIfAvailable());
}
private ReactiveMessageSenderCache reactivePulsarMessageSenderCache(
@Nullable ProducerCacheProvider producerCacheProvider) {
return (producerCacheProvider != null) ? AdaptedReactivePulsarClientFactory.createCache(producerCacheProvider)
: AdaptedReactivePulsarClientFactory.createCache();
}
@Bean
@ConditionalOnMissingBean(ReactivePulsarSenderFactory.class)
DefaultReactivePulsarSenderFactory<?> reactivePulsarSenderFactory(ReactivePulsarClient reactivePulsarClient,
ObjectProvider<ReactiveMessageSenderCache> reactiveMessageSenderCache, TopicResolver topicResolver,
ObjectProvider<ReactiveMessageSenderBuilderCustomizer<?>> customizersProvider,
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
List<ReactiveMessageSenderBuilderCustomizer<?>> customizers = new ArrayList<>();
customizers.add(this.propertiesMapper::customizeMessageSenderBuilder);
customizers.addAll(customizersProvider.orderedStream().toList());
List<ReactiveMessageSenderBuilderCustomizer<Object>> lambdaSafeCustomizers = List
.of((builder) -> applyMessageSenderBuilderCustomizers(customizers, builder));
Builder<Object> senderFactoryBuilder = DefaultReactivePulsarSenderFactory.builderFor(reactivePulsarClient)
.withDefaultConfigCustomizers(lambdaSafeCustomizers)
.withTopicResolver(topicResolver);
reactiveMessageSenderCache.ifAvailable(senderFactoryBuilder::withMessageSenderCache);
topicBuilderProvider.ifAvailable(senderFactoryBuilder::withTopicBuilder);
return senderFactoryBuilder.build();
}
@SuppressWarnings("unchecked")
private void applyMessageSenderBuilderCustomizers(List<ReactiveMessageSenderBuilderCustomizer<?>> customizers,
ReactiveMessageSenderBuilder<?> builder) {
LambdaSafe.callbacks(ReactiveMessageSenderBuilderCustomizer.class, customizers, builder)
.invoke((customizer) -> customizer.customize(builder));
}
@Bean
@ConditionalOnMissingBean(ReactivePulsarConsumerFactory.class)
DefaultReactivePulsarConsumerFactory<?> reactivePulsarConsumerFactory(
ReactivePulsarClient pulsarReactivePulsarClient,
ObjectProvider<ReactiveMessageConsumerBuilderCustomizer<?>> customizersProvider,
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
List<ReactiveMessageConsumerBuilderCustomizer<?>> customizers = new ArrayList<>();
customizers.add(this.propertiesMapper::customizeMessageConsumerBuilder);
customizers.addAll(customizersProvider.orderedStream().toList());
List<ReactiveMessageConsumerBuilderCustomizer<Object>> lambdaSafeCustomizers = List
.of((builder) -> applyMessageConsumerBuilderCustomizers(customizers, builder));
DefaultReactivePulsarConsumerFactory<?> consumerFactory = new DefaultReactivePulsarConsumerFactory<>(
pulsarReactivePulsarClient, lambdaSafeCustomizers);
topicBuilderProvider.ifAvailable(consumerFactory::setTopicBuilder);
return consumerFactory;
}
@SuppressWarnings("unchecked")
private void applyMessageConsumerBuilderCustomizers(List<ReactiveMessageConsumerBuilderCustomizer<?>> customizers,
ReactiveMessageConsumerBuilder<?> builder) {
LambdaSafe.callbacks(ReactiveMessageConsumerBuilderCustomizer.class, customizers, builder)
.invoke((customizer) -> customizer.customize(builder));
}
@Bean
@ConditionalOnMissingBean(name = "reactivePulsarListenerContainerFactory")
DefaultReactivePulsarListenerContainerFactory<?> reactivePulsarListenerContainerFactory(
ReactivePulsarConsumerFactory<Object> reactivePulsarConsumerFactory, SchemaResolver schemaResolver,
TopicResolver topicResolver, PulsarContainerFactoryCustomizers containerFactoryCustomizers) {
ReactivePulsarContainerProperties<Object> containerProperties = new ReactivePulsarContainerProperties<>();
containerProperties.setSchemaResolver(schemaResolver);
containerProperties.setTopicResolver(topicResolver);
this.propertiesMapper.customizeContainerProperties(containerProperties);
DefaultReactivePulsarListenerContainerFactory<?> containerFactory = new DefaultReactivePulsarListenerContainerFactory<>(
reactivePulsarConsumerFactory, containerProperties);
containerFactoryCustomizers.customize(containerFactory);
return containerFactory;
}
@Bean
@ConditionalOnMissingBean(ReactivePulsarReaderFactory.class)
DefaultReactivePulsarReaderFactory<?> reactivePulsarReaderFactory(ReactivePulsarClient reactivePulsarClient,
ObjectProvider<ReactiveMessageReaderBuilderCustomizer<?>> customizersProvider,
ObjectProvider<PulsarTopicBuilder> topicBuilderProvider) {
List<ReactiveMessageReaderBuilderCustomizer<?>> customizers = new ArrayList<>();
customizers.add(this.propertiesMapper::customizeMessageReaderBuilder);
customizers.addAll(customizersProvider.orderedStream().toList());
List<ReactiveMessageReaderBuilderCustomizer<Object>> lambdaSafeCustomizers = List
.of((builder) -> applyMessageReaderBuilderCustomizers(customizers, builder));
DefaultReactivePulsarReaderFactory<?> readerFactory = new DefaultReactivePulsarReaderFactory<>(
reactivePulsarClient, lambdaSafeCustomizers);
topicBuilderProvider.ifAvailable(readerFactory::setTopicBuilder);
return readerFactory;
}
@SuppressWarnings("unchecked")
private void applyMessageReaderBuilderCustomizers(List<ReactiveMessageReaderBuilderCustomizer<?>> customizers,
ReactiveMessageReaderBuilder<?> builder) {
LambdaSafe.callbacks(ReactiveMessageReaderBuilderCustomizer.class, customizers, builder)
.invoke((customizer) -> customizer.customize(builder));
}
@Bean
@ConditionalOnMissingBean
ReactivePulsarTemplate<?> pulsarReactiveTemplate(ReactivePulsarSenderFactory<?> reactivePulsarSenderFactory,
SchemaResolver schemaResolver, TopicResolver topicResolver) {
return new ReactivePulsarTemplate<>(reactivePulsarSenderFactory, schemaResolver, topicResolver);
}
@Configuration(proxyBeanMethods = false)
@EnableReactivePulsar
@ConditionalOnMissingBean(
name = PulsarAnnotationSupportBeanNames.REACTIVE_PULSAR_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)
static class EnableReactivePulsarConfiguration {
}
}
@@ -1,111 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.util.ArrayList;
import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder;
import org.springframework.boot.context.properties.PropertyMapper;
import org.springframework.pulsar.reactive.listener.ReactivePulsarContainerProperties;
/**
* Helper class used to map reactive {@link PulsarProperties} to various builder
* customizers.
*
* @author Chris Bono
* @author Phillip Webb
* @author Vedran Pavic
*/
final class PulsarReactivePropertiesMapper {
private final PulsarProperties properties;
PulsarReactivePropertiesMapper(PulsarProperties properties) {
this.properties = properties;
}
<T> void customizeMessageSenderBuilder(ReactiveMessageSenderBuilder<T> builder) {
PulsarProperties.Producer properties = this.properties.getProducer();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getName).to(builder::producerName);
map.from(properties::getTopicName).to(builder::topic);
map.from(properties::getSendTimeout).to(builder::sendTimeout);
map.from(properties::getMessageRoutingMode).to(builder::messageRoutingMode);
map.from(properties::getHashingScheme).to(builder::hashingScheme);
map.from(properties::isBatchingEnabled).to(builder::batchingEnabled);
map.from(properties::isChunkingEnabled).to(builder::chunkingEnabled);
map.from(properties::getCompressionType).to(builder::compressionType);
map.from(properties::getAccessMode).to(builder::accessMode);
}
<T> void customizeMessageConsumerBuilder(ReactiveMessageConsumerBuilder<T> builder) {
PulsarProperties.Consumer properties = this.properties.getConsumer();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getName).to(builder::consumerName);
map.from(properties::getTopics).as(ArrayList::new).to(builder::topics);
map.from(properties::getTopicsPattern).to(builder::topicsPattern);
map.from(properties::getPriorityLevel).to(builder::priorityLevel);
map.from(properties::isReadCompacted).to(builder::readCompacted);
map.from(properties::getDeadLetterPolicy).as(DeadLetterPolicyMapper::map).to(builder::deadLetterPolicy);
map.from(properties::isRetryEnable).to(builder::retryLetterTopicEnable);
customizerMessageConsumerBuilderSubscription(builder);
}
private <T> void customizerMessageConsumerBuilderSubscription(ReactiveMessageConsumerBuilder<T> builder) {
PulsarProperties.Consumer.Subscription properties = this.properties.getConsumer().getSubscription();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getName).to(builder::subscriptionName);
map.from(properties::getInitialPosition).to(builder::subscriptionInitialPosition);
map.from(properties::getMode).to(builder::subscriptionMode);
map.from(properties::getTopicsMode).to(builder::topicsPatternSubscriptionMode);
map.from(properties::getType).to(builder::subscriptionType);
}
<T> void customizeContainerProperties(ReactivePulsarContainerProperties<T> containerProperties) {
customizePulsarContainerConsumerSubscriptionProperties(containerProperties);
customizePulsarContainerListenerProperties(containerProperties);
}
private void customizePulsarContainerConsumerSubscriptionProperties(
ReactivePulsarContainerProperties<?> containerProperties) {
PulsarProperties.Consumer.Subscription properties = this.properties.getConsumer().getSubscription();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getType).to(containerProperties::setSubscriptionType);
map.from(properties::getName).to(containerProperties::setSubscriptionName);
}
private void customizePulsarContainerListenerProperties(ReactivePulsarContainerProperties<?> containerProperties) {
PulsarProperties.Listener properties = this.properties.getListener();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getSchemaType).to(containerProperties::setSchemaType);
map.from(properties::getConcurrency).to(containerProperties::setConcurrency);
}
void customizeMessageReaderBuilder(ReactiveMessageReaderBuilder<?> builder) {
PulsarProperties.Reader properties = this.properties.getReader();
PropertyMapper map = PropertyMapper.get();
map.from(properties::getName).to(builder::readerName);
map.from(properties::getTopics).to(builder::topics);
map.from(properties::getSubscriptionName).to(builder::subscriptionName);
map.from(properties::getSubscriptionRolePrefix).to(builder::generatedSubscriptionNamePrefix);
map.from(properties::isReadCompacted).to(builder::readCompacted);
}
}
@@ -1,2 +1 @@
org.springframework.boot.pulsar.autoconfigure.PulsarAutoConfiguration
org.springframework.boot.pulsar.autoconfigure.PulsarReactiveAutoConfiguration
@@ -16,24 +16,36 @@
package org.springframework.boot.pulsar.autoconfigure;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import com.github.benmanes.caffeine.cache.Caffeine;
import org.apache.pulsar.client.admin.PulsarAdminBuilder;
import org.apache.pulsar.client.api.ClientBuilder;
import org.apache.pulsar.client.api.ConsumerBuilder;
import org.apache.pulsar.client.api.ProducerBuilder;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.ReaderBuilder;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.interceptor.ProducerInterceptor;
import org.apache.pulsar.client.impl.AutoClusterFailover;
import org.apache.pulsar.common.schema.KeyValueEncodingType;
import org.apache.pulsar.common.schema.SchemaType;
import org.assertj.core.api.InstanceOfAssertFactories;
import org.assertj.core.api.ThrowingConsumer;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledForJreRange;
import org.junit.jupiter.api.condition.JRE;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentMatchers;
import org.mockito.InOrder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.AutoConfigurations;
@@ -62,7 +74,10 @@ import org.springframework.pulsar.core.DefaultPulsarReaderFactory;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.DefaultTopicResolver;
import org.springframework.pulsar.core.ProducerBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdministration;
import org.springframework.pulsar.core.PulsarClientBuilderCustomizer;
import org.springframework.pulsar.core.PulsarClientFactory;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.core.PulsarProducerFactory;
import org.springframework.pulsar.core.PulsarReaderFactory;
@@ -70,13 +85,16 @@ import org.springframework.pulsar.core.PulsarTemplate;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.ReaderBuilderCustomizer;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.function.PulsarFunctionAdministration;
import org.springframework.pulsar.listener.PulsarContainerProperties.TransactionSettings;
import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory;
import org.springframework.pulsar.transaction.PulsarAwareTransactionManager;
import org.springframework.test.util.ReflectionTestUtils;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
/**
@@ -124,8 +142,7 @@ class PulsarAutoConfigurationTests {
@Test
void autoConfiguresBeans() {
this.contextRunner.run((context) -> assertThat(context).hasSingleBean(PulsarConfiguration.class)
.hasSingleBean(PulsarConnectionDetails.class)
this.contextRunner.run((context) -> assertThat(context).hasSingleBean(PulsarConnectionDetails.class)
.hasSingleBean(DefaultPulsarClientFactory.class)
.hasSingleBean(PulsarClient.class)
.hasSingleBean(PulsarTopicBuilder.class)
@@ -150,6 +167,346 @@ class PulsarAutoConfigurationTests {
.run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class));
}
@Test
void whenHasUserDefinedConnectionDetailsBeanDoesNotAutoConfigureBean() {
PulsarConnectionDetails customConnectionDetails = mock(PulsarConnectionDetails.class);
this.contextRunner
.withBean("customPulsarConnectionDetails", PulsarConnectionDetails.class, () -> customConnectionDetails)
.run((context) -> assertThat(context).getBean(PulsarConnectionDetails.class)
.isSameAs(customConnectionDetails));
}
@Test
void whenHasUserDefinedContainerFactoryCustomizersBeanDoesNotAutoConfigureBean() {
PulsarContainerFactoryCustomizers customizers = mock(PulsarContainerFactoryCustomizers.class);
this.contextRunner
.withBean("customContainerFactoryCustomizers", PulsarContainerFactoryCustomizers.class, () -> customizers)
.run((context) -> assertThat(context).getBean(PulsarContainerFactoryCustomizers.class)
.isSameAs(customizers));
}
@Nested
class ClientTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedClientFactoryBeanDoesNotAutoConfigureBean() {
PulsarClientFactory customFactory = mock(PulsarClientFactory.class);
given(customFactory.createClient()).willReturn(mock(PulsarClient.class));
new ApplicationContextRunner().withConfiguration(AutoConfigurations.of(PulsarAutoConfiguration.class))
.withBean("customPulsarClientFactory", PulsarClientFactory.class, () -> customFactory)
.run((context) -> assertThat(context).getBean(PulsarClientFactory.class).isSameAs(customFactory));
}
@Test
void whenHasUserDefinedClientBeanDoesNotAutoConfigureBean() {
PulsarClient customClient = mock(PulsarClient.class);
new ApplicationContextRunner().withConfiguration(AutoConfigurations.of(PulsarAutoConfiguration.class))
.withBean("customPulsarClient", PulsarClient.class, () -> customClient)
.run((context) -> assertThat(context).getBean(PulsarClient.class).isSameAs(customClient));
}
@Test
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getBrokerUrl()).willReturn("connectiondetails");
this.contextRunner.withUserConfiguration(ClientTests.PulsarClientBuilderCustomizersConfig.class)
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.client.service-url=properties")
.run((context) -> {
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
Customizers<PulsarClientBuilderCustomizer, ClientBuilder> customizers = Customizers
.of(ClientBuilder.class, PulsarClientBuilderCustomizer::customize);
assertThat(customizers.fromField(clientFactory, "customizer")).callsInOrder(
ClientBuilder::serviceUrl, "connectiondetails", "fromCustomizer1", "fromCustomizer2");
});
}
@Test
void whenHasUserDefinedFailoverPropertiesAddsToClient() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getBrokerUrl()).willReturn("connectiondetails");
this.contextRunner.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.client.service-url=properties",
"spring.pulsar.client.failover.backup-clusters[0].service-url=backup-cluster-1",
"spring.pulsar.client.failover.delay=15s",
"spring.pulsar.client.failover.switch-back-delay=30s",
"spring.pulsar.client.failover.check-interval=5s",
"spring.pulsar.client.failover.backup-clusters[1].service-url=backup-cluster-2",
"spring.pulsar.client.failover.backup-clusters[1].authentication.plugin-class-name="
+ MockAuthentication.class.getName(),
"spring.pulsar.client.failover.backup-clusters[1].authentication.param.token=1234")
.run((context) -> {
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
PulsarProperties pulsarProperties = context.getBean(PulsarProperties.class);
ClientBuilder target = mock(ClientBuilder.class);
BiConsumer<PulsarClientBuilderCustomizer, ClientBuilder> customizeAction = PulsarClientBuilderCustomizer::customize;
PulsarClientBuilderCustomizer pulsarClientBuilderCustomizer = (PulsarClientBuilderCustomizer) ReflectionTestUtils
.getField(clientFactory, "customizer");
customizeAction.accept(pulsarClientBuilderCustomizer, target);
InOrder ordered = inOrder(target);
ordered.verify(target).serviceUrlProvider(ArgumentMatchers.any(AutoClusterFailover.class));
assertThat(pulsarProperties.getClient().getFailover().getDelay()).isEqualTo(Duration.ofSeconds(15));
assertThat(pulsarProperties.getClient().getFailover().getSwitchBackDelay())
.isEqualTo(Duration.ofSeconds(30));
assertThat(pulsarProperties.getClient().getFailover().getCheckInterval())
.isEqualTo(Duration.ofSeconds(5));
assertThat(pulsarProperties.getClient().getFailover().getBackupClusters().size()).isEqualTo(2);
});
}
@TestConfiguration(proxyBeanMethods = false)
static class PulsarClientBuilderCustomizersConfig {
@Bean
@Order(200)
PulsarClientBuilderCustomizer customizerFoo() {
return (builder) -> builder.serviceUrl("fromCustomizer2");
}
@Bean
@Order(100)
PulsarClientBuilderCustomizer customizerBar() {
return (builder) -> builder.serviceUrl("fromCustomizer1");
}
}
}
@Nested
class AdministrationTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class);
this.contextRunner
.withBean("customPulsarAdministration", PulsarAdministration.class, () -> pulsarAdministration)
.run((context) -> assertThat(context).getBean(PulsarAdministration.class)
.isSameAs(pulsarAdministration));
}
@Test
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getAdminUrl()).willReturn("connectiondetails");
this.contextRunner.withUserConfiguration(AdministrationTests.PulsarAdminBuilderCustomizersConfig.class)
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.admin.service-url=property")
.run((context) -> {
PulsarAdministration pulsarAdmin = context.getBean(PulsarAdministration.class);
Customizers<PulsarAdminBuilderCustomizer, PulsarAdminBuilder> customizers = Customizers
.of(PulsarAdminBuilder.class, PulsarAdminBuilderCustomizer::customize);
assertThat(customizers.fromField(pulsarAdmin, "adminCustomizers")).callsInOrder(
PulsarAdminBuilder::serviceHttpUrl, "connectiondetails", "fromCustomizer1",
"fromCustomizer2");
});
}
@TestConfiguration(proxyBeanMethods = false)
static class PulsarAdminBuilderCustomizersConfig {
@Bean
@Order(200)
PulsarAdminBuilderCustomizer customizerFoo() {
return (builder) -> builder.serviceHttpUrl("fromCustomizer2");
}
@Bean
@Order(100)
PulsarAdminBuilderCustomizer customizerBar() {
return (builder) -> builder.serviceHttpUrl("fromCustomizer1");
}
}
}
@Nested
class SchemaResolverTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
SchemaResolver schemaResolver = mock(SchemaResolver.class);
this.contextRunner.withBean("customSchemaResolver", SchemaResolver.class, () -> schemaResolver)
.run((context) -> assertThat(context).getBean(SchemaResolver.class).isSameAs(schemaResolver));
}
@Test
void whenHasUserDefinedSchemaResolverCustomizer() {
SchemaResolverCustomizer<DefaultSchemaResolver> customizer = (schemaResolver) -> schemaResolver
.addCustomSchemaMapping(TestRecord.class, Schema.STRING);
this.contextRunner.withBean("schemaResolverCustomizer", SchemaResolverCustomizer.class, () -> customizer)
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, Schema.STRING)));
}
@Test
void whenHasDefaultsTypeMappingForPrimitiveAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=STRING");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, Schema.STRING)));
}
@Test
void whenHasDefaultsTypeMappingForStructAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=JSON");
Schema<?> expectedSchema = Schema.JSON(TestRecord.class);
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, expectedSchema)));
}
@Test
void whenHasDefaultsTypeMappingForKeyValueAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=key-value");
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.message-key-type=java.lang.String");
Schema<?> expectedSchema = Schema.KeyValue(Schema.STRING, Schema.JSON(TestRecord.class),
KeyValueEncodingType.INLINE);
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, expectedSchema)));
}
private ThrowingConsumer<DefaultSchemaResolver> customSchemaMappingOf(Class<?> messageType,
Schema<?> expectedSchema) {
return (resolver) -> assertThat(resolver.getCustomSchemaMapping(messageType))
.hasValueSatisfying(schemaEqualTo(expectedSchema));
}
private Consumer<Schema<?>> schemaEqualTo(Schema<?> expected) {
return (actual) -> assertThat(actual.getSchemaInfo()).isEqualTo(expected.getSchemaInfo());
}
}
@Nested
class TopicResolverTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
TopicResolver topicResolver = mock(TopicResolver.class);
this.contextRunner.withBean("customTopicResolver", TopicResolver.class, () -> topicResolver)
.run((context) -> assertThat(context).getBean(TopicResolver.class).isSameAs(topicResolver));
}
@Test
void whenHasDefaultsTypeMappingAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].topic-name=foo-topic");
properties.add("spring.pulsar.defaults.type-mappings[1].message-type=java.lang.String");
properties.add("spring.pulsar.defaults.type-mappings[1].topic-name=string-topic");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(TopicResolver.class)
.asInstanceOf(InstanceOfAssertFactories.type(DefaultTopicResolver.class))
.satisfies((resolver) -> {
assertThat(resolver.getCustomTopicMapping(TestRecord.class)).hasValue("foo-topic");
assertThat(resolver.getCustomTopicMapping(String.class)).hasValue("string-topic");
}));
}
}
@Nested
class TopicBuilderTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class);
this.contextRunner.withBean("customPulsarTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder)
.run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class).isSameAs(topicBuilder));
}
@Test
void whenHasDefaultsTopicDisabledPropertyDoesNotCreateBean() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class));
}
@Test
void whenHasDefaultsTenantAndNamespaceAppliedToTopicBuilder() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.topic.tenant=my-tenant");
properties.add("spring.pulsar.defaults.topic.namespace=my-namespace");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class)
.asInstanceOf(InstanceOfAssertFactories.type(PulsarTopicBuilder.class))
.satisfies((topicBuilder) -> {
assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultTenant", "my-tenant");
assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultNamespace", "my-namespace");
}));
}
@Test
void beanHasScopePrototype() {
this.contextRunner.run((context) -> assertThat(context.getBean(PulsarTopicBuilder.class))
.isNotSameAs(context.getBean(PulsarTopicBuilder.class)));
}
}
@Nested
class FunctionAdministrationTests {
private final ApplicationContextRunner contextRunner = PulsarAutoConfigurationTests.this.contextRunner;
@Test
void whenNoPropertiesAddsFunctionAdministrationBean() {
this.contextRunner.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.hasFieldOrPropertyWithValue("failFast", Boolean.TRUE)
.hasFieldOrPropertyWithValue("propagateFailures", Boolean.TRUE)
.hasFieldOrPropertyWithValue("propagateStopFailures", Boolean.FALSE)
.hasNoNullFieldsOrProperties() // ensures object providers set
.extracting("pulsarAdministration")
.isSameAs(context.getBean(PulsarAdministration.class)));
}
@Test
void whenHasFunctionPropertiesAppliesPropertiesToBean() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.function.fail-fast=false");
properties.add("spring.pulsar.function.propagate-failures=false");
properties.add("spring.pulsar.function.propagate-stop-failures=true");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.hasFieldOrPropertyWithValue("failFast", Boolean.FALSE)
.hasFieldOrPropertyWithValue("propagateFailures", Boolean.FALSE)
.hasFieldOrPropertyWithValue("propagateStopFailures", Boolean.TRUE));
}
@Test
void whenHasFunctionDisabledPropertyDoesNotCreateBean() {
this.contextRunner.withPropertyValues("spring.pulsar.function.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(PulsarFunctionAdministration.class));
}
@Test
void whenHasCustomFunctionAdministrationBean() {
PulsarFunctionAdministration functionAdministration = mock(PulsarFunctionAdministration.class);
this.contextRunner.withBean(PulsarFunctionAdministration.class, () -> functionAdministration)
.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.isSameAs(functionAdministration));
}
}
@Nested
class ProducerFactoryTests {
@@ -597,14 +954,6 @@ class PulsarAutoConfigurationTests {
@TestConfiguration(proxyBeanMethods = false)
static class ListenerContainerFactoryCustomizersConfig {
@Bean
@Order(50)
PulsarContainerFactoryCustomizer<DefaultReactivePulsarListenerContainerFactory<?>> customizerIgnored() {
return (containerFactory) -> {
throw new IllegalStateException("should-not-have-matched");
};
}
@Bean
@Order(200)
PulsarContainerFactoryCustomizer<ConcurrentPulsarListenerContainerFactory<?>> customizerFoo() {
@@ -723,14 +1072,6 @@ class PulsarAutoConfigurationTests {
@TestConfiguration(proxyBeanMethods = false)
static class ReaderContainerFactoryCustomizersConfig {
@Bean
@Order(50)
PulsarContainerFactoryCustomizer<DefaultReactivePulsarListenerContainerFactory<?>> customizerIgnored() {
return (containerFactory) -> {
throw new IllegalStateException("should-not-have-matched");
};
}
@Bean
@Order(200)
PulsarContainerFactoryCustomizer<DefaultPulsarReaderContainerFactory<?>> customizerFoo() {
@@ -787,4 +1128,10 @@ class PulsarAutoConfigurationTests {
}
record TestRecord() {
private static final String CLASS_NAME = TestRecord.class.getName();
}
}
@@ -1,421 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
import org.apache.pulsar.client.admin.PulsarAdminBuilder;
import org.apache.pulsar.client.api.ClientBuilder;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.impl.AutoClusterFailover;
import org.apache.pulsar.common.schema.KeyValueEncodingType;
import org.assertj.core.api.InstanceOfAssertFactories;
import org.assertj.core.api.ThrowingConsumer;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentMatchers;
import org.mockito.InOrder;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.core.annotation.Order;
import org.springframework.pulsar.core.DefaultPulsarClientFactory;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.DefaultTopicResolver;
import org.springframework.pulsar.core.PulsarAdminBuilderCustomizer;
import org.springframework.pulsar.core.PulsarAdministration;
import org.springframework.pulsar.core.PulsarClientBuilderCustomizer;
import org.springframework.pulsar.core.PulsarClientFactory;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.function.PulsarFunctionAdministration;
import org.springframework.test.util.ReflectionTestUtils;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
/**
* Tests for {@link PulsarConfiguration}.
*
* @author Chris Bono
* @author Alexander Preuß
* @author Soby Chacko
* @author Phillip Webb
* @author Swamy Mavuri
*/
class PulsarConfigurationTests {
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(PulsarConfiguration.class))
.withBean(PulsarClient.class, () -> mock(PulsarClient.class));
@Test
void whenHasUserDefinedConnectionDetailsBeanDoesNotAutoConfigureBean() {
PulsarConnectionDetails customConnectionDetails = mock(PulsarConnectionDetails.class);
this.contextRunner
.withBean("customPulsarConnectionDetails", PulsarConnectionDetails.class, () -> customConnectionDetails)
.run((context) -> assertThat(context).getBean(PulsarConnectionDetails.class)
.isSameAs(customConnectionDetails));
}
@Test
void whenHasUserDefinedContainerFactoryCustomizersBeanDoesNotAutoConfigureBean() {
PulsarContainerFactoryCustomizers customizers = mock(PulsarContainerFactoryCustomizers.class);
this.contextRunner
.withBean("customContainerFactoryCustomizers", PulsarContainerFactoryCustomizers.class, () -> customizers)
.run((context) -> assertThat(context).getBean(PulsarContainerFactoryCustomizers.class)
.isSameAs(customizers));
}
@Nested
class ClientTests {
@Test
void whenHasUserDefinedClientFactoryBeanDoesNotAutoConfigureBean() {
PulsarClientFactory customFactory = mock(PulsarClientFactory.class);
new ApplicationContextRunner().withConfiguration(AutoConfigurations.of(PulsarConfiguration.class))
.withBean("customPulsarClientFactory", PulsarClientFactory.class, () -> customFactory)
.run((context) -> assertThat(context).getBean(PulsarClientFactory.class).isSameAs(customFactory));
}
@Test
void whenHasUserDefinedClientBeanDoesNotAutoConfigureBean() {
PulsarClient customClient = mock(PulsarClient.class);
new ApplicationContextRunner().withConfiguration(AutoConfigurations.of(PulsarConfiguration.class))
.withBean("customPulsarClient", PulsarClient.class, () -> customClient)
.run((context) -> assertThat(context).getBean(PulsarClient.class).isSameAs(customClient));
}
@Test
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getBrokerUrl()).willReturn("connectiondetails");
PulsarConfigurationTests.this.contextRunner
.withUserConfiguration(PulsarClientBuilderCustomizersConfig.class)
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.client.service-url=properties")
.run((context) -> {
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
Customizers<PulsarClientBuilderCustomizer, ClientBuilder> customizers = Customizers
.of(ClientBuilder.class, PulsarClientBuilderCustomizer::customize);
assertThat(customizers.fromField(clientFactory, "customizer")).callsInOrder(
ClientBuilder::serviceUrl, "connectiondetails", "fromCustomizer1", "fromCustomizer2");
});
}
@Test
void whenHasUserDefinedFailoverPropertiesAddsToClient() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getBrokerUrl()).willReturn("connectiondetails");
PulsarConfigurationTests.this.contextRunner.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.client.service-url=properties",
"spring.pulsar.client.failover.backup-clusters[0].service-url=backup-cluster-1",
"spring.pulsar.client.failover.delay=15s",
"spring.pulsar.client.failover.switch-back-delay=30s",
"spring.pulsar.client.failover.check-interval=5s",
"spring.pulsar.client.failover.backup-clusters[1].service-url=backup-cluster-2",
"spring.pulsar.client.failover.backup-clusters[1].authentication.plugin-class-name="
+ MockAuthentication.class.getName(),
"spring.pulsar.client.failover.backup-clusters[1].authentication.param.token=1234")
.run((context) -> {
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
PulsarProperties pulsarProperties = context.getBean(PulsarProperties.class);
ClientBuilder target = mock(ClientBuilder.class);
BiConsumer<PulsarClientBuilderCustomizer, ClientBuilder> customizeAction = PulsarClientBuilderCustomizer::customize;
PulsarClientBuilderCustomizer pulsarClientBuilderCustomizer = (PulsarClientBuilderCustomizer) ReflectionTestUtils
.getField(clientFactory, "customizer");
customizeAction.accept(pulsarClientBuilderCustomizer, target);
InOrder ordered = inOrder(target);
ordered.verify(target).serviceUrlProvider(ArgumentMatchers.any(AutoClusterFailover.class));
assertThat(pulsarProperties.getClient().getFailover().getDelay()).isEqualTo(Duration.ofSeconds(15));
assertThat(pulsarProperties.getClient().getFailover().getSwitchBackDelay())
.isEqualTo(Duration.ofSeconds(30));
assertThat(pulsarProperties.getClient().getFailover().getCheckInterval())
.isEqualTo(Duration.ofSeconds(5));
assertThat(pulsarProperties.getClient().getFailover().getBackupClusters().size()).isEqualTo(2);
});
}
@TestConfiguration(proxyBeanMethods = false)
static class PulsarClientBuilderCustomizersConfig {
@Bean
@Order(200)
PulsarClientBuilderCustomizer customizerFoo() {
return (builder) -> builder.serviceUrl("fromCustomizer2");
}
@Bean
@Order(100)
PulsarClientBuilderCustomizer customizerBar() {
return (builder) -> builder.serviceUrl("fromCustomizer1");
}
}
}
@Nested
class AdministrationTests {
private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class);
this.contextRunner
.withBean("customPulsarAdministration", PulsarAdministration.class, () -> pulsarAdministration)
.run((context) -> assertThat(context).getBean(PulsarAdministration.class)
.isSameAs(pulsarAdministration));
}
@Test
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
given(connectionDetails.getAdminUrl()).willReturn("connectiondetails");
this.contextRunner.withUserConfiguration(PulsarAdminBuilderCustomizersConfig.class)
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
.withPropertyValues("spring.pulsar.admin.service-url=property")
.run((context) -> {
PulsarAdministration pulsarAdmin = context.getBean(PulsarAdministration.class);
Customizers<PulsarAdminBuilderCustomizer, PulsarAdminBuilder> customizers = Customizers
.of(PulsarAdminBuilder.class, PulsarAdminBuilderCustomizer::customize);
assertThat(customizers.fromField(pulsarAdmin, "adminCustomizers")).callsInOrder(
PulsarAdminBuilder::serviceHttpUrl, "connectiondetails", "fromCustomizer1",
"fromCustomizer2");
});
}
@TestConfiguration(proxyBeanMethods = false)
static class PulsarAdminBuilderCustomizersConfig {
@Bean
@Order(200)
PulsarAdminBuilderCustomizer customizerFoo() {
return (builder) -> builder.serviceHttpUrl("fromCustomizer2");
}
@Bean
@Order(100)
PulsarAdminBuilderCustomizer customizerBar() {
return (builder) -> builder.serviceHttpUrl("fromCustomizer1");
}
}
}
@Nested
class SchemaResolverTests {
private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
SchemaResolver schemaResolver = mock(SchemaResolver.class);
this.contextRunner.withBean("customSchemaResolver", SchemaResolver.class, () -> schemaResolver)
.run((context) -> assertThat(context).getBean(SchemaResolver.class).isSameAs(schemaResolver));
}
@Test
void whenHasUserDefinedSchemaResolverCustomizer() {
SchemaResolverCustomizer<DefaultSchemaResolver> customizer = (schemaResolver) -> schemaResolver
.addCustomSchemaMapping(TestRecord.class, Schema.STRING);
this.contextRunner.withBean("schemaResolverCustomizer", SchemaResolverCustomizer.class, () -> customizer)
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, Schema.STRING)));
}
@Test
void whenHasDefaultsTypeMappingForPrimitiveAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=STRING");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, Schema.STRING)));
}
@Test
void whenHasDefaultsTypeMappingForStructAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=JSON");
Schema<?> expectedSchema = Schema.JSON(TestRecord.class);
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, expectedSchema)));
}
@Test
void whenHasDefaultsTypeMappingForKeyValueAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=key-value");
properties.add("spring.pulsar.defaults.type-mappings[0].schema-info.message-key-type=java.lang.String");
Schema<?> expectedSchema = Schema.KeyValue(Schema.STRING, Schema.JSON(TestRecord.class),
KeyValueEncodingType.INLINE);
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(DefaultSchemaResolver.class)
.satisfies(customSchemaMappingOf(TestRecord.class, expectedSchema)));
}
private ThrowingConsumer<DefaultSchemaResolver> customSchemaMappingOf(Class<?> messageType,
Schema<?> expectedSchema) {
return (resolver) -> assertThat(resolver.getCustomSchemaMapping(messageType))
.hasValueSatisfying(schemaEqualTo(expectedSchema));
}
private Consumer<Schema<?>> schemaEqualTo(Schema<?> expected) {
return (actual) -> assertThat(actual.getSchemaInfo()).isEqualTo(expected.getSchemaInfo());
}
}
@Nested
class TopicResolverTests {
private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
TopicResolver topicResolver = mock(TopicResolver.class);
this.contextRunner.withBean("customTopicResolver", TopicResolver.class, () -> topicResolver)
.run((context) -> assertThat(context).getBean(TopicResolver.class).isSameAs(topicResolver));
}
@Test
void whenHasDefaultsTypeMappingAddsToSchemaResolver() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.type-mappings[0].message-type=" + TestRecord.CLASS_NAME);
properties.add("spring.pulsar.defaults.type-mappings[0].topic-name=foo-topic");
properties.add("spring.pulsar.defaults.type-mappings[1].message-type=java.lang.String");
properties.add("spring.pulsar.defaults.type-mappings[1].topic-name=string-topic");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(TopicResolver.class)
.asInstanceOf(InstanceOfAssertFactories.type(DefaultTopicResolver.class))
.satisfies((resolver) -> {
assertThat(resolver.getCustomTopicMapping(TestRecord.class)).hasValue("foo-topic");
assertThat(resolver.getCustomTopicMapping(String.class)).hasValue("string-topic");
}));
}
}
@Nested
class TopicBuilderTests {
private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class);
this.contextRunner.withBean("customPulsarTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder)
.run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class).isSameAs(topicBuilder));
}
@Test
void whenHasDefaultsTopicDisabledPropertyDoesNotCreateBean() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class));
}
@Test
void whenHasDefaultsTenantAndNamespaceAppliedToTopicBuilder() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.defaults.topic.tenant=my-tenant");
properties.add("spring.pulsar.defaults.topic.namespace=my-namespace");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(PulsarTopicBuilder.class)
.asInstanceOf(InstanceOfAssertFactories.type(PulsarTopicBuilder.class))
.satisfies((topicBuilder) -> {
assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultTenant", "my-tenant");
assertThat(topicBuilder).hasFieldOrPropertyWithValue("defaultNamespace", "my-namespace");
}));
}
@Test
void beanHasScopePrototype() {
this.contextRunner.run((context) -> assertThat(context.getBean(PulsarTopicBuilder.class))
.isNotSameAs(context.getBean(PulsarTopicBuilder.class)));
}
}
@Nested
class FunctionAdministrationTests {
private final ApplicationContextRunner contextRunner = PulsarConfigurationTests.this.contextRunner;
@Test
void whenNoPropertiesAddsFunctionAdministrationBean() {
this.contextRunner.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.hasFieldOrPropertyWithValue("failFast", Boolean.TRUE)
.hasFieldOrPropertyWithValue("propagateFailures", Boolean.TRUE)
.hasFieldOrPropertyWithValue("propagateStopFailures", Boolean.FALSE)
.hasNoNullFieldsOrProperties() // ensures object providers set
.extracting("pulsarAdministration")
.isSameAs(context.getBean(PulsarAdministration.class)));
}
@Test
void whenHasFunctionPropertiesAppliesPropertiesToBean() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.function.fail-fast=false");
properties.add("spring.pulsar.function.propagate-failures=false");
properties.add("spring.pulsar.function.propagate-stop-failures=true");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new))
.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.hasFieldOrPropertyWithValue("failFast", Boolean.FALSE)
.hasFieldOrPropertyWithValue("propagateFailures", Boolean.FALSE)
.hasFieldOrPropertyWithValue("propagateStopFailures", Boolean.TRUE));
}
@Test
void whenHasFunctionDisabledPropertyDoesNotCreateBean() {
this.contextRunner.withPropertyValues("spring.pulsar.function.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(PulsarFunctionAdministration.class));
}
@Test
void whenHasCustomFunctionAdministrationBean() {
PulsarFunctionAdministration functionAdministration = mock(PulsarFunctionAdministration.class);
this.contextRunner.withBean(PulsarFunctionAdministration.class, () -> functionAdministration)
.run((context) -> assertThat(context).getBean(PulsarFunctionAdministration.class)
.isSameAs(functionAdministration));
}
}
record TestRecord() {
private static final String CLASS_NAME = TestRecord.class.getName();
}
}
@@ -26,10 +26,8 @@ import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactor
import org.springframework.pulsar.config.DefaultPulsarReaderContainerFactory;
import org.springframework.pulsar.config.ListenerContainerFactory;
import org.springframework.pulsar.config.PulsarContainerFactory;
import org.springframework.pulsar.config.PulsarListenerContainerFactory;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.listener.PulsarContainerProperties;
import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.then;
@@ -74,12 +72,12 @@ class PulsarContainerFactoryCustomizersTests {
assertThat(list.get(1).getCount()).isZero();
assertThat(list.get(2).getCount()).isZero();
customizers.customize(mock(ConcurrentPulsarListenerContainerFactory.class));
customizers.customize(mock(ListenerContainerFactory.class));
assertThat(list.get(0).getCount()).isEqualTo(2);
assertThat(list.get(1).getCount()).isOne();
assertThat(list.get(2).getCount()).isOne();
assertThat(list.get(2).getCount()).isZero();
customizers.customize(mock(DefaultReactivePulsarListenerContainerFactory.class));
customizers.customize(mock(ConcurrentPulsarListenerContainerFactory.class));
assertThat(list.get(0).getCount()).isEqualTo(3);
assertThat(list.get(1).getCount()).isEqualTo(2);
assertThat(list.get(2).getCount()).isOne();
@@ -101,7 +99,7 @@ class PulsarContainerFactoryCustomizersTests {
}
/**
* Test customizer that will match all {@link PulsarListenerContainerFactory}.
* Test customizer that will match all {@link PulsarContainerFactory}.
*
* @param <T> the container factory type
*/
@@ -121,10 +119,7 @@ class PulsarContainerFactoryCustomizersTests {
}
/**
* Test customizer that will match both
* {@link ConcurrentPulsarListenerContainerFactory} and
* {@link DefaultReactivePulsarListenerContainerFactory} as they both extend
* {@link ListenerContainerFactory}.
* Test customizer that will match all {@link ListenerContainerFactory}.
*/
static class TestPulsarListenersContainerFactoryCustomizer extends TestCustomizer<ListenerContainerFactory<?, ?>> {
@@ -1,555 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Supplier;
import com.github.benmanes.caffeine.cache.Caffeine;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.common.schema.SchemaType;
import org.apache.pulsar.reactive.client.adapter.ProducerCacheProvider;
import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderCache;
import org.apache.pulsar.reactive.client.api.ReactivePulsarClient;
import org.apache.pulsar.reactive.client.producercache.CaffeineShadedProducerCacheProvider;
import org.assertj.core.api.AbstractObjectAssert;
import org.assertj.core.api.InstanceOfAssertFactories;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.test.context.FilteredClassLoader;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.boot.test.context.assertj.AssertableApplicationContext;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.core.annotation.Order;
import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.DefaultTopicResolver;
import org.springframework.pulsar.core.PulsarAdministration;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.pulsar.core.TopicResolver;
import org.springframework.pulsar.reactive.config.DefaultReactivePulsarListenerContainerFactory;
import org.springframework.pulsar.reactive.config.ReactivePulsarListenerContainerFactory;
import org.springframework.pulsar.reactive.config.ReactivePulsarListenerEndpointRegistry;
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarBootstrapConfiguration;
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListenerAnnotationBeanPostProcessor;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarConsumerFactory;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarReaderFactory;
import org.springframework.pulsar.reactive.core.DefaultReactivePulsarSenderFactory;
import org.springframework.pulsar.reactive.core.ReactiveMessageConsumerBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactiveMessageReaderBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactiveMessageSenderBuilderCustomizer;
import org.springframework.pulsar.reactive.core.ReactivePulsarConsumerFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarReaderFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarSenderFactory;
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate;
import org.springframework.pulsar.reactive.listener.ReactivePulsarContainerProperties;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
/**
* Tests for {@link PulsarReactiveAutoConfiguration}.
*
* @author Chris Bono
* @author Christophe Bornet
* @author Phillip Webb
*/
class PulsarReactiveAutoConfigurationTests {
private static final String INTERNAL_PULSAR_LISTENER_ANNOTATION_PROCESSOR = "org.springframework.pulsar.config.internalReactivePulsarListenerAnnotationProcessor";
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(PulsarReactiveAutoConfiguration.class))
.withBean(PulsarClient.class, () -> mock(PulsarClient.class));
@Test
void whenPulsarNotOnClasspathAutoConfigurationIsSkipped() {
new ApplicationContextRunner().withConfiguration(AutoConfigurations.of(PulsarReactiveAutoConfiguration.class))
.withClassLoader(new FilteredClassLoader(PulsarClient.class))
.run((context) -> assertThat(context).doesNotHaveBean(PulsarReactiveAutoConfiguration.class));
}
@Test
void whenReactivePulsarNotOnClasspathAutoConfigurationIsSkipped() {
this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarClient.class))
.run((context) -> assertThat(context).doesNotHaveBean(PulsarReactiveAutoConfiguration.class));
}
@Test
void whenReactiveSpringPulsarNotOnClasspathAutoConfigurationIsSkipped() {
this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarTemplate.class))
.run((context) -> assertThat(context).doesNotHaveBean(PulsarReactiveAutoConfiguration.class));
}
@Test
void whenCustomPulsarListenerAnnotationProcessorDefinedAutoConfigurationIsSkipped() {
this.contextRunner.withBean(INTERNAL_PULSAR_LISTENER_ANNOTATION_PROCESSOR, String.class, () -> "bean")
.run((context) -> assertThat(context).doesNotHaveBean(ReactivePulsarBootstrapConfiguration.class));
}
@Test
void autoConfiguresBeans() {
this.contextRunner.run((context) -> assertThat(context).hasSingleBean(PulsarConfiguration.class)
.hasSingleBean(PulsarClient.class)
.hasSingleBean(PulsarTopicBuilder.class)
.hasSingleBean(PulsarAdministration.class)
.hasSingleBean(DefaultSchemaResolver.class)
.hasSingleBean(DefaultTopicResolver.class)
.hasSingleBean(ReactivePulsarClient.class)
.hasSingleBean(CaffeineShadedProducerCacheProvider.class)
.hasSingleBean(ReactiveMessageSenderCache.class)
.hasSingleBean(DefaultReactivePulsarSenderFactory.class)
.hasSingleBean(ReactivePulsarTemplate.class)
.hasSingleBean(DefaultReactivePulsarConsumerFactory.class)
.hasSingleBean(DefaultReactivePulsarListenerContainerFactory.class)
.hasSingleBean(ReactivePulsarListenerAnnotationBeanPostProcessor.class)
.hasSingleBean(ReactivePulsarListenerEndpointRegistry.class));
}
@Test
void topicDefaultsCanBeDisabled() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(PulsarTopicBuilder.class));
}
@Test
@SuppressWarnings("rawtypes")
void injectsExpectedBeansIntoReactivePulsarClient() {
this.contextRunner.run((context) -> {
PulsarClient pulsarClient = context.getBean(PulsarClient.class);
assertThat(context).hasNotFailed()
.getBean(ReactivePulsarClient.class)
.extracting("reactivePulsarResourceAdapter")
.extracting("pulsarClientSupplier", InstanceOfAssertFactories.type(Supplier.class))
.extracting(Supplier::get)
.isSameAs(pulsarClient);
});
}
@ParameterizedTest
@ValueSource(classes = { ReactivePulsarClient.class, ProducerCacheProvider.class, ReactiveMessageSenderCache.class,
ReactivePulsarSenderFactory.class, ReactivePulsarConsumerFactory.class, ReactivePulsarReaderFactory.class,
ReactivePulsarTemplate.class })
<T> void whenHasUserDefinedBeanDoesNotAutoConfigureBean(Class<T> beanClass) {
T bean = mock(beanClass);
this.contextRunner.withBean(beanClass.getName(), beanClass, () -> bean)
.run((context) -> assertThat(context).getBean(beanClass).isSameAs(bean));
}
@Nested
class SenderFactoryTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
void injectsExpectedBeans() {
ReactivePulsarClient client = mock(ReactivePulsarClient.class);
ReactiveMessageSenderCache cache = mock(ReactiveMessageSenderCache.class);
this.contextRunner.withPropertyValues("spring.pulsar.producer.topic-name=test-topic")
.withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client)
.withBean("customReactiveMessageSenderCache", ReactiveMessageSenderCache.class, () -> cache)
.run((context) -> {
DefaultReactivePulsarSenderFactory<?> senderFactory = context
.getBean(DefaultReactivePulsarSenderFactory.class);
assertThat(senderFactory)
.extracting("reactivePulsarClient", InstanceOfAssertFactories.type(ReactivePulsarClient.class))
.isSameAs(client);
assertThat(senderFactory)
.extracting("reactiveMessageSenderCache",
InstanceOfAssertFactories.type(ReactiveMessageSenderCache.class))
.isSameAs(cache);
assertThat(senderFactory)
.extracting("topicResolver", InstanceOfAssertFactories.type(TopicResolver.class))
.isSameAs(context.getBean(TopicResolver.class));
assertThat(senderFactory).extracting("topicBuilder").isNotNull();
});
}
@Test
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat((DefaultReactivePulsarSenderFactory<?>) context
.getBean(DefaultReactivePulsarSenderFactory.class)).extracting("topicBuilder").isNull());
}
@Test
void injectsExpectedBeansIntoReactiveMessageSenderCache() {
ProducerCacheProvider provider = mock(ProducerCacheProvider.class);
this.contextRunner.withBean("customProducerCacheProvider", ProducerCacheProvider.class, () -> provider)
.run((context) -> assertThat(context).getBean(ReactiveMessageSenderCache.class)
.extracting("cacheProvider", InstanceOfAssertFactories.type(ProducerCacheProvider.class))
.isSameAs(provider));
}
@Test
<T> void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
this.contextRunner.withPropertyValues("spring.pulsar.producer.name=fromPropsCustomizer")
.withUserConfiguration(ReactiveMessageSenderBuilderCustomizerConfig.class)
.run((context) -> {
DefaultReactivePulsarSenderFactory<?> producerFactory = context
.getBean(DefaultReactivePulsarSenderFactory.class);
Customizers<ReactiveMessageSenderBuilderCustomizer<T>, ReactiveMessageSenderBuilder<T>> customizers = Customizers
.of(ReactiveMessageSenderBuilder.class, ReactiveMessageSenderBuilderCustomizer::customize);
assertThat(customizers.fromField(producerFactory, "defaultConfigCustomizers")).callsInOrder(
ReactiveMessageSenderBuilder::producerName, "fromPropsCustomizer", "fromCustomizer1",
"fromCustomizer2");
});
}
@TestConfiguration(proxyBeanMethods = false)
static class ReactiveMessageSenderBuilderCustomizerConfig {
@Bean
@Order(200)
ReactiveMessageSenderBuilderCustomizer<?> customizerFoo() {
return (builder) -> builder.producerName("fromCustomizer2");
}
@Bean
@Order(100)
ReactiveMessageSenderBuilderCustomizer<?> customizerBar() {
return (builder) -> builder.producerName("fromCustomizer1");
}
}
}
@Nested
class TemplateTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
@SuppressWarnings("rawtypes")
void injectsExpectedBeans() {
ReactivePulsarSenderFactory senderFactory = mock(ReactivePulsarSenderFactory.class);
SchemaResolver schemaResolver = mock(SchemaResolver.class);
this.contextRunner
.withBean("customReactivePulsarSenderFactory", ReactivePulsarSenderFactory.class, () -> senderFactory)
.withBean("schemaResolver", SchemaResolver.class, () -> schemaResolver)
.run((context) -> assertThat(context).getBean(ReactivePulsarTemplate.class).satisfies((template) -> {
assertThat(template).extracting("reactiveMessageSenderFactory").isSameAs(senderFactory);
assertThat(template).extracting("schemaResolver").isSameAs(schemaResolver);
}));
}
}
@Nested
class ConsumerFactoryTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
void injectsExpectedBeans() {
ReactivePulsarClient client = mock(ReactivePulsarClient.class);
PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class);
this.contextRunner.withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client)
.withBean("customTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder)
.run((context) -> {
ReactivePulsarConsumerFactory<?> consumerFactory = context
.getBean(DefaultReactivePulsarConsumerFactory.class);
assertThat(consumerFactory)
.extracting("reactivePulsarClient", InstanceOfAssertFactories.type(ReactivePulsarClient.class))
.isSameAs(client);
assertThat(consumerFactory)
.extracting("topicBuilder", InstanceOfAssertFactories.type(PulsarTopicBuilder.class))
.isSameAs(topicBuilder);
});
}
@Test
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat(
(ReactivePulsarConsumerFactory<?>) context.getBean(DefaultReactivePulsarConsumerFactory.class))
.extracting("topicBuilder")
.isNull());
}
@Test
<T> void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
this.contextRunner.withPropertyValues("spring.pulsar.consumer.name=fromPropsCustomizer")
.withUserConfiguration(ReactiveMessageConsumerBuilderCustomizerConfig.class)
.run((context) -> {
DefaultReactivePulsarConsumerFactory<?> consumerFactory = context
.getBean(DefaultReactivePulsarConsumerFactory.class);
Customizers<ReactiveMessageConsumerBuilderCustomizer<T>, ReactiveMessageConsumerBuilder<T>> customizers = Customizers
.of(ReactiveMessageConsumerBuilder.class, ReactiveMessageConsumerBuilderCustomizer::customize);
assertThat(customizers.fromField(consumerFactory, "defaultConfigCustomizers")).callsInOrder(
ReactiveMessageConsumerBuilder::consumerName, "fromPropsCustomizer", "fromCustomizer1",
"fromCustomizer2");
});
}
@TestConfiguration(proxyBeanMethods = false)
static class ReactiveMessageConsumerBuilderCustomizerConfig {
@Bean
@Order(200)
ReactiveMessageConsumerBuilderCustomizer<?> customizerFoo() {
return (builder) -> builder.consumerName("fromCustomizer2");
}
@Bean
@Order(100)
ReactiveMessageConsumerBuilderCustomizer<?> customizerBar() {
return (builder) -> builder.consumerName("fromCustomizer1");
}
}
}
@Nested
class ListenerTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
void whenHasUserDefinedBeanDoesNotAutoConfigureBean() {
ReactivePulsarListenerContainerFactory<?> listenerContainerFactory = mock(
ReactivePulsarListenerContainerFactory.class);
this.contextRunner
.withBean("reactivePulsarListenerContainerFactory", ReactivePulsarListenerContainerFactory.class,
() -> listenerContainerFactory)
.run((context) -> assertThat(context).getBean(ReactivePulsarListenerContainerFactory.class)
.isSameAs(listenerContainerFactory));
}
@Test
void whenHasUserDefinedReactivePulsarListenerAnnotationBeanPostProcessorDoesNotAutoConfigureBean() {
ReactivePulsarListenerAnnotationBeanPostProcessor<?> listenerAnnotationBeanPostProcessor = mock(
ReactivePulsarListenerAnnotationBeanPostProcessor.class);
this.contextRunner.withBean(INTERNAL_PULSAR_LISTENER_ANNOTATION_PROCESSOR,
ReactivePulsarListenerAnnotationBeanPostProcessor.class, () -> listenerAnnotationBeanPostProcessor)
.run((context) -> assertThat(context).getBean(ReactivePulsarListenerAnnotationBeanPostProcessor.class)
.isSameAs(listenerAnnotationBeanPostProcessor));
}
@Test
void whenHasCustomProperties() {
List<String> properties = new ArrayList<>();
properties.add("spring.pulsar.listener.schema-type=avro");
this.contextRunner.withPropertyValues(properties.toArray(String[]::new)).run((context) -> {
DefaultReactivePulsarListenerContainerFactory<?> factory = context
.getBean(DefaultReactivePulsarListenerContainerFactory.class);
assertThat(factory.getContainerProperties().getSchemaType()).isEqualTo(SchemaType.AVRO);
});
}
@Test
void injectsExpectedBeans() {
ReactivePulsarConsumerFactory<?> consumerFactory = mock(ReactivePulsarConsumerFactory.class);
SchemaResolver schemaResolver = mock(SchemaResolver.class);
this.contextRunner
.withBean("customReactivePulsarConsumerFactory", ReactivePulsarConsumerFactory.class,
() -> consumerFactory)
.withBean("schemaResolver", SchemaResolver.class, () -> schemaResolver)
.run((context) -> {
DefaultReactivePulsarListenerContainerFactory<?> containerFactory = context
.getBean(DefaultReactivePulsarListenerContainerFactory.class);
assertThat(containerFactory).extracting("consumerFactory").isSameAs(consumerFactory);
assertThat(containerFactory)
.extracting(DefaultReactivePulsarListenerContainerFactory::getContainerProperties)
.extracting(ReactivePulsarContainerProperties::getSchemaResolver)
.isSameAs(schemaResolver);
});
}
@Test
void whenHasUserDefinedFactoryCustomizersAppliesInCorrectOrder() {
this.contextRunner.withUserConfiguration(ListenerContainerFactoryCustomizersConfig.class)
.run((context) -> assertThat(context).getBean(DefaultReactivePulsarListenerContainerFactory.class)
.hasFieldOrPropertyWithValue("containerProperties.subscriptionName", ":bar:foo"));
}
@TestConfiguration(proxyBeanMethods = false)
static class ListenerContainerFactoryCustomizersConfig {
@Bean
@Order(50)
PulsarContainerFactoryCustomizer<ConcurrentPulsarListenerContainerFactory<?>> customizerIgnored() {
return (containerFactory) -> {
throw new IllegalStateException("should-not-have-matched");
};
}
@Bean
@Order(200)
PulsarContainerFactoryCustomizer<DefaultReactivePulsarListenerContainerFactory<?>> customizerFoo() {
return (containerFactory) -> appendToSubscriptionName(containerFactory, ":foo");
}
@Bean
@Order(100)
PulsarContainerFactoryCustomizer<DefaultReactivePulsarListenerContainerFactory<?>> customizerBar() {
return (containerFactory) -> appendToSubscriptionName(containerFactory, ":bar");
}
private void appendToSubscriptionName(DefaultReactivePulsarListenerContainerFactory<?> containerFactory,
String valueToAppend) {
String subscriptionName = containerFactory.getContainerProperties().getSubscriptionName();
String updatedValue = (subscriptionName != null) ? subscriptionName + valueToAppend : valueToAppend;
containerFactory.getContainerProperties().setSubscriptionName(updatedValue);
}
}
}
@Nested
class ReaderFactoryTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
void injectsExpectedBeans() {
ReactivePulsarClient client = mock(ReactivePulsarClient.class);
PulsarTopicBuilder topicBuilder = mock(PulsarTopicBuilder.class);
this.contextRunner.withPropertyValues("spring.pulsar.reader.name=test-reader")
.withBean("customReactivePulsarClient", ReactivePulsarClient.class, () -> client)
.withBean("customPulsarTopicBuilder", PulsarTopicBuilder.class, () -> topicBuilder)
.run((context) -> {
DefaultReactivePulsarReaderFactory<?> readerFactory = context
.getBean(DefaultReactivePulsarReaderFactory.class);
assertThat(readerFactory)
.extracting("reactivePulsarClient", InstanceOfAssertFactories.type(ReactivePulsarClient.class))
.isSameAs(client);
assertThat(readerFactory)
.extracting("topicBuilder", InstanceOfAssertFactories.type(PulsarTopicBuilder.class))
.isSameAs(topicBuilder);
});
}
@Test
void hasNoTopicBuilderWhenTopicDefaultsAreDisabled() {
this.contextRunner.withPropertyValues("spring.pulsar.defaults.topic.enabled=false")
.run((context) -> assertThat((DefaultReactivePulsarReaderFactory<?>) context
.getBean(DefaultReactivePulsarReaderFactory.class)).extracting("topicBuilder").isNull());
}
@Test
<T> void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
this.contextRunner.withPropertyValues("spring.pulsar.reader.name=fromPropsCustomizer")
.withUserConfiguration(ReactiveMessageReaderBuilderCustomizerConfig.class)
.run((context) -> {
DefaultReactivePulsarReaderFactory<?> readerFactory = context
.getBean(DefaultReactivePulsarReaderFactory.class);
Customizers<ReactiveMessageReaderBuilderCustomizer<T>, ReactiveMessageReaderBuilder<T>> customizers = Customizers
.of(ReactiveMessageReaderBuilder.class, ReactiveMessageReaderBuilderCustomizer::customize);
assertThat(customizers.fromField(readerFactory, "defaultConfigCustomizers")).callsInOrder(
ReactiveMessageReaderBuilder::readerName, "fromPropsCustomizer", "fromCustomizer1",
"fromCustomizer2");
});
}
@TestConfiguration(proxyBeanMethods = false)
static class ReactiveMessageReaderBuilderCustomizerConfig {
@Bean
@Order(200)
ReactiveMessageReaderBuilderCustomizer<?> customizerFoo() {
return (builder) -> builder.readerName("fromCustomizer2");
}
@Bean
@Order(100)
ReactiveMessageReaderBuilderCustomizer<?> customizerBar() {
return (builder) -> builder.readerName("fromCustomizer1");
}
}
}
@Nested
class SenderCacheAutoConfigurationTests {
private final ApplicationContextRunner contextRunner = PulsarReactiveAutoConfigurationTests.this.contextRunner;
@Test
void whenNoPropertiesEnablesCaching() {
this.contextRunner.run(this::assertCaffeineProducerCacheProvider);
}
@Test
void whenCachingEnabledEnablesCaching() {
this.contextRunner.withPropertyValues("spring.pulsar.producer.cache.enabled=true")
.run(this::assertCaffeineProducerCacheProvider);
}
@Test
void whenCachingDisabledDoesNotEnableCaching() {
this.contextRunner.withPropertyValues("spring.pulsar.producer.cache.enabled=false")
.run((context) -> assertThat(context).doesNotHaveBean(ProducerCacheProvider.class)
.doesNotHaveBean(ReactiveMessageSenderCache.class));
}
@Test
void whenCachingEnabledAndCaffeineNotOnClasspathStillUsesCaffeine() {
// The reactive client shades Caffeine - it should still be used
this.contextRunner.withClassLoader(new FilteredClassLoader(Caffeine.class))
.withPropertyValues("spring.pulsar.producer.cache.enabled=true")
.run(this::assertCaffeineProducerCacheProvider);
}
@Test
void whenCachingEnabledAndNoCacheProviderAvailable() {
// The reactive client uses a shaded caffeine cache provider as its internal
// cache
this.contextRunner.withClassLoader(new FilteredClassLoader(CaffeineShadedProducerCacheProvider.class))
.withPropertyValues("spring.pulsar.producer.cache.enabled=true")
.run((context) -> assertThat(context).doesNotHaveBean(ProducerCacheProvider.class)
.getBean(ReactiveMessageSenderCache.class)
.extracting("cacheProvider")
.isExactlyInstanceOf(CaffeineShadedProducerCacheProvider.class));
}
@Test
void whenCustomCachingPropertiesCreatesConfiguredBean() {
this.contextRunner
.withPropertyValues("spring.pulsar.producer.cache.expire-after-access=100s",
"spring.pulsar.producer.cache.maximum-size=5150",
"spring.pulsar.producer.cache.initial-capacity=200")
.run((context) -> assertCaffeineProducerCacheProvider(context).extracting("cache.cache")
.hasFieldOrPropertyWithValue("expiresAfterAccessNanos", Duration.ofSeconds(100).toNanos())
.hasFieldOrPropertyWithValue("maximum", 5150L));
}
private AbstractObjectAssert<?, ProducerCacheProvider> assertCaffeineProducerCacheProvider(
AssertableApplicationContext context) {
return assertThat(context).hasSingleBean(ReactiveMessageSenderCache.class)
.getBean(ProducerCacheProvider.class)
.isExactlyInstanceOf(CaffeineShadedProducerCacheProvider.class);
}
}
}
@@ -1,151 +0,0 @@
/*
* Copyright 2012-present 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 org.springframework.boot.pulsar.autoconfigure;
import java.time.Duration;
import java.util.List;
import java.util.regex.Pattern;
import org.apache.pulsar.client.api.CompressionType;
import org.apache.pulsar.client.api.DeadLetterPolicy;
import org.apache.pulsar.client.api.HashingScheme;
import org.apache.pulsar.client.api.MessageRoutingMode;
import org.apache.pulsar.client.api.ProducerAccessMode;
import org.apache.pulsar.client.api.RegexSubscriptionMode;
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
import org.apache.pulsar.client.api.SubscriptionMode;
import org.apache.pulsar.client.api.SubscriptionType;
import org.apache.pulsar.common.schema.SchemaType;
import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageReaderBuilder;
import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderBuilder;
import org.junit.jupiter.api.Test;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Consumer;
import org.springframework.boot.pulsar.autoconfigure.PulsarProperties.Consumer.Subscription;
import org.springframework.pulsar.reactive.listener.ReactivePulsarContainerProperties;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.mock;
/**
* Tests for {@link PulsarReactivePropertiesMapper}.
*
* @author Chris Bono
* @author Phillip Webb
* @author Vedran Pavic
*/
class PulsarReactivePropertiesMapperTests {
@Test
@SuppressWarnings("unchecked")
void customizeMessageSenderBuilder() {
PulsarProperties properties = new PulsarProperties();
properties.getProducer().setName("name");
properties.getProducer().setTopicName("topicname");
properties.getProducer().setSendTimeout(Duration.ofSeconds(1));
properties.getProducer().setMessageRoutingMode(MessageRoutingMode.RoundRobinPartition);
properties.getProducer().setHashingScheme(HashingScheme.JavaStringHash);
properties.getProducer().setBatchingEnabled(false);
properties.getProducer().setChunkingEnabled(true);
properties.getProducer().setCompressionType(CompressionType.SNAPPY);
properties.getProducer().setAccessMode(ProducerAccessMode.Exclusive);
ReactiveMessageSenderBuilder<Object> builder = mock(ReactiveMessageSenderBuilder.class);
new PulsarReactivePropertiesMapper(properties).customizeMessageSenderBuilder(builder);
then(builder).should().producerName("name");
then(builder).should().topic("topicname");
then(builder).should().sendTimeout(Duration.ofSeconds(1));
then(builder).should().messageRoutingMode(MessageRoutingMode.RoundRobinPartition);
then(builder).should().hashingScheme(HashingScheme.JavaStringHash);
then(builder).should().batchingEnabled(false);
then(builder).should().chunkingEnabled(true);
then(builder).should().compressionType(CompressionType.SNAPPY);
then(builder).should().accessMode(ProducerAccessMode.Exclusive);
}
@Test
@SuppressWarnings("unchecked")
void customizeMessageConsumerBuilder() {
PulsarProperties properties = new PulsarProperties();
List<String> topics = List.of("mytopic");
Pattern topisPattern = Pattern.compile("my-pattern");
properties.getConsumer().setName("name");
properties.getConsumer().setTopics(topics);
properties.getConsumer().setTopicsPattern(topisPattern);
properties.getConsumer().setPriorityLevel(123);
properties.getConsumer().setReadCompacted(true);
Consumer.DeadLetterPolicy deadLetterPolicy = new Consumer.DeadLetterPolicy();
deadLetterPolicy.setDeadLetterTopic("my-dlt");
deadLetterPolicy.setMaxRedeliverCount(1);
properties.getConsumer().setDeadLetterPolicy(deadLetterPolicy);
properties.getConsumer().setRetryEnable(false);
Subscription subscriptionProperties = properties.getConsumer().getSubscription();
subscriptionProperties.setName("subname");
subscriptionProperties.setInitialPosition(SubscriptionInitialPosition.Earliest);
subscriptionProperties.setMode(SubscriptionMode.NonDurable);
subscriptionProperties.setTopicsMode(RegexSubscriptionMode.NonPersistentOnly);
subscriptionProperties.setType(SubscriptionType.Key_Shared);
ReactiveMessageConsumerBuilder<Object> builder = mock(ReactiveMessageConsumerBuilder.class);
new PulsarReactivePropertiesMapper(properties).customizeMessageConsumerBuilder(builder);
then(builder).should().consumerName("name");
then(builder).should().topics(topics);
then(builder).should().topicsPattern(topisPattern);
then(builder).should().priorityLevel(123);
then(builder).should().readCompacted(true);
then(builder).should().deadLetterPolicy(new DeadLetterPolicy(1, null, "my-dlt", null, null, null));
then(builder).should().retryLetterTopicEnable(false);
then(builder).should().subscriptionName("subname");
then(builder).should().subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
then(builder).should().subscriptionMode(SubscriptionMode.NonDurable);
then(builder).should().topicsPatternSubscriptionMode(RegexSubscriptionMode.NonPersistentOnly);
then(builder).should().subscriptionType(SubscriptionType.Key_Shared);
}
@Test
void customizeContainerProperties() {
PulsarProperties properties = new PulsarProperties();
properties.getConsumer().getSubscription().setType(SubscriptionType.Shared);
properties.getConsumer().getSubscription().setName("my-subscription");
properties.getListener().setSchemaType(SchemaType.AVRO);
properties.getListener().setConcurrency(10);
ReactivePulsarContainerProperties<Object> containerProperties = new ReactivePulsarContainerProperties<>();
new PulsarReactivePropertiesMapper(properties).customizeContainerProperties(containerProperties);
assertThat(containerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Shared);
assertThat(containerProperties.getSubscriptionName()).isEqualTo("my-subscription");
assertThat(containerProperties.getSchemaType()).isEqualTo(SchemaType.AVRO);
assertThat(containerProperties.getConcurrency()).isEqualTo(10);
}
@Test
@SuppressWarnings("unchecked")
void customizeMessageReaderBuilder() {
List<String> topics = List.of("my-topic");
PulsarProperties properties = new PulsarProperties();
properties.getReader().setName("name");
properties.getReader().setTopics(topics);
properties.getReader().setSubscriptionName("subname");
properties.getReader().setSubscriptionRolePrefix("srp");
ReactiveMessageReaderBuilder<Object> builder = mock(ReactiveMessageReaderBuilder.class);
new PulsarReactivePropertiesMapper(properties).customizeMessageReaderBuilder(builder);
then(builder).should().readerName("name");
then(builder).should().topics(topics);
then(builder).should().subscriptionName("subname");
then(builder).should().generatedSubscriptionNamePrefix("srp");
}
}
@@ -1796,15 +1796,6 @@ bom {
releaseNotes("https://pulsar.apache.org/release-notes/versioned/pulsar-{version}")
}
}
library("Pulsar Reactive", "0.7.0") {
group("org.apache.pulsar") {
bom("pulsar-client-reactive-bom")
}
links {
site("https://github.com/apache/pulsar-client-reactive")
releaseNotes("https://github.com/apache/pulsar-client-reactive/releases/tag/v{version}")
}
}
library("Quartz", "2.5.0") {
group("org.quartz-scheduler") {
modules = [
@@ -2207,8 +2198,6 @@ bom {
"spring-boot-starter-opentelemetry-test",
"spring-boot-starter-pulsar",
"spring-boot-starter-pulsar-test",
"spring-boot-starter-pulsar-reactive",
"spring-boot-starter-pulsar-reactive-test",
"spring-boot-starter-quartz",
"spring-boot-starter-quartz-test",
"spring-boot-starter-r2dbc",
-2
View File
@@ -310,8 +310,6 @@ include "starter:spring-boot-starter-opentelemetry-test"
include "starter:spring-boot-starter-parent"
include "starter:spring-boot-starter-pulsar"
include "starter:spring-boot-starter-pulsar-test"
include "starter:spring-boot-starter-pulsar-reactive"
include "starter:spring-boot-starter-pulsar-reactive-test"
include "starter:spring-boot-starter-quartz"
include "starter:spring-boot-starter-quartz-test"
include "starter:spring-boot-starter-r2dbc"
@@ -23,7 +23,6 @@ description = "Spring Boot Pulsar smoke test"
dependencies {
implementation(project(":starter:spring-boot-starter-pulsar"))
implementation(project(":starter:spring-boot-starter-pulsar-reactive"))
dockerTestImplementation(project(":starter:spring-boot-starter-test"))
dockerTestImplementation(project(":test-support:spring-boot-docker-test-support"))
@@ -19,10 +19,7 @@ package smoketest.pulsar;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.IntStream;
import org.awaitility.Awaitility;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.testcontainers.junit.jupiter.Container;
@@ -34,59 +31,27 @@ import org.springframework.boot.test.system.CapturedOutput;
import org.springframework.boot.test.system.OutputCaptureExtension;
import org.springframework.boot.testcontainers.service.connection.ServiceConnection;
import org.springframework.boot.testsupport.container.TestImage;
import org.springframework.test.context.ActiveProfiles;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.waitAtMost;
@Testcontainers(disabledWithoutDocker = true)
@ExtendWith(OutputCaptureExtension.class)
@SpringBootTest
class SamplePulsarApplicationTests {
@Container
@ServiceConnection
static final PulsarContainer pulsar = TestImage.container(PulsarContainer.class);
abstract class PulsarApplication {
private final String type;
PulsarApplication(String type) {
this.type = type;
@Test
void appProducesAndConsumesMessages(CapturedOutput output) {
List<String> expectedOutput = new ArrayList<>();
for (int i = 0; i < 10; i++) {
expectedOutput.add("++++++PRODUCE:(%s)------".formatted(i));
expectedOutput.add("++++++CONSUME:(%s)------".formatted(i));
}
@Test
void appProducesAndConsumesMessages(CapturedOutput output) {
List<String> expectedOutput = new ArrayList<>();
IntStream.range(0, 10).forEachOrdered((i) -> {
expectedOutput.add("++++++PRODUCE %s:(%s)------".formatted(this.type, i));
expectedOutput.add("++++++CONSUME %s:(%s)------".formatted(this.type, i));
});
Awaitility.waitAtMost(Duration.ofSeconds(30))
.untilAsserted(() -> assertThat(output).contains(expectedOutput));
}
}
@Nested
@SpringBootTest
@ActiveProfiles("smoketest-pulsar-imperative")
class ImperativePulsarApplication extends PulsarApplication {
ImperativePulsarApplication() {
super("IMPERATIVE");
}
}
@Nested
@SpringBootTest
@ActiveProfiles("smoketest-pulsar-reactive")
class ReactivePulsarApplication extends PulsarApplication {
ReactivePulsarApplication() {
super("REACTIVE");
}
waitAtMost(Duration.ofSeconds(30)).untilAsserted(() -> assertThat(output).contains(expectedOutput));
}
}
@@ -1,64 +0,0 @@
/*
* Copyright 2012-present 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 smoketest.pulsar;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.pulsar.reactive.client.api.MessageSpec;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.pulsar.core.PulsarTopic;
import org.springframework.pulsar.core.PulsarTopicBuilder;
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener;
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate;
@Configuration(proxyBeanMethods = false)
@Profile("smoketest-pulsar-reactive")
class ReactiveAppConfig {
private static final Log logger = LogFactory.getLog(ReactiveAppConfig.class);
private static final String TOPIC = "pulsar-reactive-smoke-test-topic";
@Bean
PulsarTopic pulsarTestTopic() {
return new PulsarTopicBuilder().name(TOPIC).numberOfPartitions(1).build();
}
@Bean
ApplicationRunner sendMessagesToPulsarTopic(ReactivePulsarTemplate<SampleMessage> template) {
return (args) -> Flux.range(0, 10)
.map((i) -> new SampleMessage(i, "message:" + i))
.map(MessageSpec::of)
.as((msgs) -> template.send(TOPIC, msgs))
.doOnNext((sendResult) -> logger
.info("++++++PRODUCE REACTIVE:(" + sendResult.getMessageSpec().getValue().id() + ")------"))
.subscribe();
}
@ReactivePulsarListener(topics = TOPIC)
Mono<Void> consumeMessagesFromPulsarTopic(SampleMessage msg) {
logger.info("++++++CONSUME REACTIVE:(" + msg.id() + ")------");
return Mono.empty();
}
}
@@ -22,17 +22,15 @@ import org.apache.commons.logging.LogFactory;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.pulsar.annotation.PulsarListener;
import org.springframework.pulsar.core.PulsarTemplate;
import org.springframework.pulsar.core.PulsarTopic;
import org.springframework.pulsar.core.PulsarTopicBuilder;
@Configuration(proxyBeanMethods = false)
@Profile("smoketest-pulsar-imperative")
class ImperativeAppConfig {
class SamplePulsarApplicationConfig {
private static final Log logger = LogFactory.getLog(ImperativeAppConfig.class);
private static final Log logger = LogFactory.getLog(SamplePulsarApplicationConfig.class);
private static final String TOPIC = "pulsar-smoke-test-topic";
@@ -46,14 +44,14 @@ class ImperativeAppConfig {
return (args) -> {
for (int i = 0; i < 10; i++) {
template.send(TOPIC, new SampleMessage(i, "message:" + i));
logger.info("++++++PRODUCE IMPERATIVE:(" + i + ")------");
logger.info("++++++PRODUCE:(" + i + ")------");
}
};
}
@PulsarListener(topics = TOPIC)
void consumeMessagesFromPulsarTopic(SampleMessage msg) {
logger.info("++++++CONSUME IMPERATIVE:(" + msg.id() + ")------");
logger.info("++++++CONSUME:(" + msg.id() + ")------");
}
}
@@ -1,32 +0,0 @@
/*
* Copyright 2012-present 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.
*/
plugins {
id "org.springframework.boot.starter"
}
description = "Starter for testing Spring for Apache Pulsar Reactive"
dependencies {
api(project(":starter:spring-boot-starter-pulsar-reactive"))
api(project(":starter:spring-boot-starter-test"))
}
checkRuntimeClasspathForConflicts {
ignore { name -> name.startsWith("org/bouncycastle/") ||
name.matches("^org/apache/pulsar/.*/package-info.class\$") ||
name.equals("findbugsExclude.xml") }
}
@@ -1,35 +0,0 @@
/*
* Copyright 2012-present 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.
*/
plugins {
id "org.springframework.boot.starter"
}
description = "Starter for using Spring for Apache Pulsar Reactive"
dependencies {
api(project(":starter:spring-boot-starter"))
api(project(":module:spring-boot-pulsar"))
api(project(":module:spring-boot-reactor"))
api("org.springframework.pulsar:spring-pulsar-reactive")
}
checkRuntimeClasspathForConflicts {
ignore { name -> name.startsWith("org/bouncycastle/") ||
name.matches("^org/apache/pulsar/.*/package-info.class\$") ||
name.equals("findbugsExclude.xml") }
}