Rename ReactivePulsarSenderTemplate to ReactivePulsarTemplate
This commit is contained in:
committed by
Chris Bono
parent
7e91f78fa1
commit
82815d9dc5
@@ -38,7 +38,7 @@ import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.pulsar.annotation.PulsarListener;
|
||||
import org.springframework.pulsar.core.reactive.DefaultReactivePulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
|
||||
@@ -59,7 +59,7 @@ public class ReactiveSpringPulsarBootApp {
|
||||
private final Logger logger = LoggerFactory.getLogger(this.getClass());
|
||||
|
||||
@Autowired
|
||||
private ReactivePulsarSenderTemplate<Foo> reactivePulsarTemplate;
|
||||
private ReactivePulsarTemplate<Foo> reactivePulsarTemplate;
|
||||
|
||||
// TODO remove this once the auto-config is available
|
||||
@Bean
|
||||
@@ -103,7 +103,7 @@ public class ReactiveSpringPulsarBootApp {
|
||||
private final Logger logger = LoggerFactory.getLogger(this.getClass());
|
||||
|
||||
@Bean
|
||||
ApplicationRunner sendSimple(ReactivePulsarSenderTemplate<String> reactivePulsarTemplate) {
|
||||
ApplicationRunner sendSimple(ReactivePulsarTemplate<String> reactivePulsarTemplate) {
|
||||
return args -> reactivePulsarTemplate
|
||||
.send("sample-reactive-topic2", Flux.range(0, 10).map((i) -> "msg-from-sendSimple-" + i))
|
||||
.subscribe();
|
||||
|
||||
@@ -37,7 +37,7 @@ import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFacto
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate;
|
||||
|
||||
import com.github.benmanes.caffeine.cache.Caffeine;
|
||||
|
||||
@@ -48,7 +48,7 @@ import com.github.benmanes.caffeine.cache.Caffeine;
|
||||
* @author Christophe Bornet
|
||||
*/
|
||||
@AutoConfiguration(after = PulsarAutoConfiguration.class)
|
||||
@ConditionalOnClass({ ReactivePulsarSenderTemplate.class, ReactivePulsarClient.class })
|
||||
@ConditionalOnClass({ ReactivePulsarTemplate.class, ReactivePulsarClient.class })
|
||||
@EnableConfigurationProperties(PulsarReactiveProperties.class)
|
||||
public class PulsarReactiveAutoConfiguration {
|
||||
|
||||
@@ -111,9 +111,9 @@ public class PulsarReactiveAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public ReactivePulsarSenderTemplate<?> pulsarReactiveSenderTemplate(
|
||||
public ReactivePulsarTemplate<?> pulsarReactiveTemplate(
|
||||
ReactivePulsarSenderFactory<?> reactivePulsarSenderFactory) {
|
||||
return new ReactivePulsarSenderTemplate<>(reactivePulsarSenderFactory);
|
||||
return new ReactivePulsarTemplate<>(reactivePulsarSenderFactory);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -51,7 +51,7 @@ import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFacto
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate;
|
||||
import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate;
|
||||
|
||||
/**
|
||||
* Autoconfiguration tests for {@link PulsarReactiveAutoConfiguration}.
|
||||
@@ -72,23 +72,23 @@ class PulsarReactiveAutoConfigurationTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
void autoConfigurationSkippedWhenReactivePulsarSenderTemplateNotOnClasspath() {
|
||||
this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarSenderTemplate.class)).run(
|
||||
void autoConfigurationSkippedWhenReactivePulsarTemplateNotOnClasspath() {
|
||||
this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarTemplate.class)).run(
|
||||
(context) -> assertThat(context).hasNotFailed().doesNotHaveBean(PulsarReactiveAutoConfiguration.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
void defaultBeansAreAutoConfigured() {
|
||||
this.contextRunner.run((context) -> assertThat(context).hasNotFailed()
|
||||
.hasSingleBean(ReactivePulsarSenderTemplate.class).hasSingleBean(ReactivePulsarClient.class)
|
||||
.hasSingleBean(ReactivePulsarTemplate.class).hasSingleBean(ReactivePulsarClient.class)
|
||||
.hasSingleBean(ProducerCacheProvider.class).hasSingleBean(ReactiveMessageSenderCache.class)
|
||||
.hasSingleBean(ReactivePulsarSenderFactory.class).getBean(ReactivePulsarSenderTemplate.class));
|
||||
.hasSingleBean(ReactivePulsarSenderFactory.class).getBean(ReactivePulsarTemplate.class));
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@ValueSource(classes = { ReactivePulsarClient.class, ProducerCacheProvider.class, ReactiveMessageSenderCache.class,
|
||||
ReactivePulsarSenderFactory.class, ReactivePulsarConsumerFactory.class, ReactivePulsarReaderFactory.class,
|
||||
ReactivePulsarSenderTemplate.class })
|
||||
ReactivePulsarTemplate.class })
|
||||
<T> void customBeanIsRespected(Class<T> beanClass) {
|
||||
T bean = mock(beanClass);
|
||||
this.contextRunner.withBean(beanClass.getName(), beanClass, () -> bean)
|
||||
@@ -100,7 +100,7 @@ class PulsarReactiveAutoConfigurationTests {
|
||||
ReactivePulsarSenderFactory<?> senderFactory = mock(ReactivePulsarSenderFactory.class);
|
||||
this.contextRunner
|
||||
.withBean("customReactivePulsarSenderFactory", ReactivePulsarSenderFactory.class, () -> senderFactory)
|
||||
.run((context -> assertThat(context).hasNotFailed().getBean(ReactivePulsarSenderTemplate.class)
|
||||
.run((context -> assertThat(context).hasNotFailed().getBean(ReactivePulsarTemplate.class)
|
||||
.extracting("reactiveMessageSenderFactory")
|
||||
.asInstanceOf(InstanceOfAssertFactories.type(ReactivePulsarSenderFactory.class))
|
||||
.isSameAs(senderFactory)));
|
||||
|
||||
@@ -28,7 +28,7 @@ import reactor.core.publisher.Mono;
|
||||
* @param <T> the message payload type
|
||||
* @author Christophe Bornet
|
||||
*/
|
||||
public interface ReactivePulsarSenderOperations<T> {
|
||||
public interface ReactivePulsarOperations<T> {
|
||||
|
||||
/**
|
||||
* Sends a message to the default topic in a reactive manner.
|
||||
@@ -74,7 +74,7 @@ public interface ReactivePulsarSenderOperations<T> {
|
||||
|
||||
/**
|
||||
* Builder that can be used to configure and send a message. Provides more options
|
||||
* than the send methods provided by {@link ReactivePulsarSenderOperations}.
|
||||
* than the send methods provided by {@link ReactivePulsarOperations}.
|
||||
*
|
||||
* @param <T> the message payload type
|
||||
*/
|
||||
@@ -37,7 +37,7 @@ import reactor.core.publisher.Mono;
|
||||
* @param <T> the message payload type
|
||||
* @author Christophe Bornet
|
||||
*/
|
||||
public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOperations<T> {
|
||||
public class ReactivePulsarTemplate<T> implements ReactivePulsarOperations<T> {
|
||||
|
||||
private final LogAccessor logger = new LogAccessor(this.getClass());
|
||||
|
||||
@@ -50,7 +50,7 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
|
||||
* @param reactiveMessageSenderFactory the factory used to create the backing Pulsar
|
||||
* reactive senders
|
||||
*/
|
||||
public ReactivePulsarSenderTemplate(ReactivePulsarSenderFactory<T> reactiveMessageSenderFactory) {
|
||||
public ReactivePulsarTemplate(ReactivePulsarSenderFactory<T> reactiveMessageSenderFactory) {
|
||||
this.reactiveMessageSenderFactory = reactiveMessageSenderFactory;
|
||||
}
|
||||
|
||||
@@ -141,7 +141,7 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
|
||||
|
||||
public static class SendMessageBuilderImpl<T> implements SendMessageBuilder<T> {
|
||||
|
||||
private final ReactivePulsarSenderTemplate<T> template;
|
||||
private final ReactivePulsarTemplate<T> template;
|
||||
|
||||
private final T message;
|
||||
|
||||
@@ -151,7 +151,7 @@ public class ReactivePulsarSenderTemplate<T> implements ReactivePulsarSenderOper
|
||||
|
||||
private ReactiveMessageSenderBuilderCustomizer<T> senderCustomizer;
|
||||
|
||||
SendMessageBuilderImpl(ReactivePulsarSenderTemplate<T> template, T message) {
|
||||
SendMessageBuilderImpl(ReactivePulsarTemplate<T> template, T message) {
|
||||
this.template = template;
|
||||
this.message = message;
|
||||
}
|
||||
@@ -45,7 +45,7 @@ import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
/**
|
||||
* Tests for {@link ReactivePulsarSenderTemplate}.
|
||||
* Tests for {@link ReactivePulsarTemplate}.
|
||||
*
|
||||
* @author Christophe Bornet
|
||||
*/
|
||||
@@ -62,7 +62,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
|
||||
senderSpec.setTopicName(topic);
|
||||
ReactivePulsarSenderFactory<Foo> producerFactory = new DefaultReactivePulsarSenderFactory<>(client,
|
||||
senderSpec, null);
|
||||
ReactivePulsarSenderTemplate<Foo> pulsarTemplate = new ReactivePulsarSenderTemplate<>(producerFactory);
|
||||
ReactivePulsarTemplate<Foo> pulsarTemplate = new ReactivePulsarTemplate<>(producerFactory);
|
||||
pulsarTemplate.setSchema(Schema.JSON(Foo.class));
|
||||
|
||||
List<Foo> foos = new ArrayList<>();
|
||||
@@ -104,7 +104,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
|
||||
}
|
||||
ReactivePulsarSenderFactory<String> senderFactory = new DefaultReactivePulsarSenderFactory<>(client,
|
||||
senderSpec, null);
|
||||
ReactivePulsarSenderTemplate<String> pulsarTemplate = new ReactivePulsarSenderTemplate<>(senderFactory);
|
||||
ReactivePulsarTemplate<String> pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory);
|
||||
Mono<MessageId> sendResponse;
|
||||
if (testArgs.useTemplateSchema) {
|
||||
pulsarTemplate.setSchema(Schema.STRING);
|
||||
@@ -114,7 +114,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport {
|
||||
: pulsarTemplate.send(msgPayload);
|
||||
}
|
||||
else {
|
||||
ReactivePulsarSenderTemplate.SendMessageBuilderImpl<String> messageBuilder = pulsarTemplate
|
||||
ReactivePulsarTemplate.SendMessageBuilderImpl<String> messageBuilder = pulsarTemplate
|
||||
.newMessage(msgPayload);
|
||||
if (testArgs.useSpecificTopic) {
|
||||
messageBuilder = messageBuilder.withTopic(topic);
|
||||
|
||||
Reference in New Issue
Block a user