From cf782afb789f80c5779f0f310d7fb80e784ee3f5 Mon Sep 17 00:00:00 2001 From: Vedran Pavic <1149230+vpavic@users.noreply.github.com> Date: Sat, 16 Dec 2023 04:39:39 +0100 Subject: [PATCH] Deprecate (Reactive)PulsarListenerEndpointAdapter (#481) This commit deprecates PulsarListenerEndpointAdapter and ReactivePulsarListenerEndpointAdapter in favor of default methods on ListenerEndpoint and its sub-interfaces, which makes it a bit easier to provide custom ListenerEndpoint implementations. --- ...eactivePulsarListenerContainerFactory.java | 2 +- .../ReactivePulsarListenerEndpoint.java | 9 ++++- ...ReactivePulsarListenerEndpointAdapter.java | 2 + ...currentPulsarListenerContainerFactory.java | 2 +- .../pulsar/config/ListenerEndpoint.java | 37 ++++++++++++++----- .../pulsar/config/PulsarListenerEndpoint.java | 13 +++++-- .../config/PulsarListenerEndpointAdapter.java | 2 + 7 files changed, 51 insertions(+), 16 deletions(-) diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/DefaultReactivePulsarListenerContainerFactory.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/DefaultReactivePulsarListenerContainerFactory.java index e0eeb62a..aafe914a 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/DefaultReactivePulsarListenerContainerFactory.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/DefaultReactivePulsarListenerContainerFactory.java @@ -156,7 +156,7 @@ public class DefaultReactivePulsarListenerContainerFactory implements Reactiv @Override public DefaultReactivePulsarMessageListenerContainer createContainer(String... topics) { - ReactivePulsarListenerEndpoint endpoint = new ReactivePulsarListenerEndpointAdapter<>() { + ReactivePulsarListenerEndpoint endpoint = new ReactivePulsarListenerEndpoint<>() { @Override public List getTopics() { diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpoint.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpoint.java index f7f442d3..65e71dcb 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpoint.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpoint.java @@ -29,12 +29,17 @@ import org.springframework.pulsar.reactive.listener.ReactivePulsarMessageListene * @param Message payload type. * @author Christophe Bornet * @author Chris Bono + * @author Vedran Pavic */ public interface ReactivePulsarListenerEndpoint extends ListenerEndpoint> { - boolean isFluxListener(); + default boolean isFluxListener() { + return false; + } @Nullable - Boolean getUseKeyOrderedProcessing(); + default Boolean getUseKeyOrderedProcessing() { + return null; + } } diff --git a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpointAdapter.java b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpointAdapter.java index 9f2f28f7..e10860c9 100644 --- a/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpointAdapter.java +++ b/spring-pulsar-reactive/src/main/java/org/springframework/pulsar/reactive/config/ReactivePulsarListenerEndpointAdapter.java @@ -31,7 +31,9 @@ import org.springframework.pulsar.support.MessageConverter; * * @param Message payload type. * @author Christophe Bornet + * @deprecated for removal in favor of {@link ReactivePulsarListenerEndpoint} */ +@Deprecated(forRemoval = true) public class ReactivePulsarListenerEndpointAdapter implements ReactivePulsarListenerEndpoint { @Override diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java index 2e8c22b7..aad3e6b5 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java @@ -54,7 +54,7 @@ public class ConcurrentPulsarListenerContainerFactory @Override public ConcurrentPulsarMessageListenerContainer createContainer(String... topics) { - PulsarListenerEndpoint endpoint = new PulsarListenerEndpointAdapter() { + PulsarListenerEndpoint endpoint = new PulsarListenerEndpoint() { @Override public Collection getTopics() { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ListenerEndpoint.java index a409115b..72e7fa06 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ListenerEndpoint.java @@ -17,6 +17,7 @@ package org.springframework.pulsar.config; import java.util.Collection; +import java.util.Collections; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.schema.SchemaType; @@ -32,6 +33,7 @@ import org.springframework.pulsar.support.MessageConverter; * * @param Message listener container type. * @author Christophe Bornet + * @author Vedran Pavic */ public interface ListenerEndpoint { @@ -42,53 +44,69 @@ public interface ListenerEndpoint { * @see ListenerContainerFactory#createListenerContainer */ @Nullable - String getId(); + default String getId() { + return null; + } /** * Return the subscription name for this endpoint's container. * @return the subscription name. */ @Nullable - String getSubscriptionName(); + default String getSubscriptionName() { + return null; + } /** * Return the subscription type for this endpoint's container. * @return the subscription type. */ @Nullable - SubscriptionType getSubscriptionType(); + default SubscriptionType getSubscriptionType() { + return SubscriptionType.Exclusive; + } /** * Return the topics for this endpoint's container. * @return the topics. */ - Collection getTopics(); + default Collection getTopics() { + return Collections.emptyList(); + } /** * Return the topic pattern for this endpoint's container. * @return the topic pattern. */ - String getTopicPattern(); + default String getTopicPattern() { + return null; + } /** * Return the autoStartup for this endpoint's container. * @return the autoStartup. */ @Nullable - Boolean getAutoStartup(); + default Boolean getAutoStartup() { + return null; + } /** * Return the schema type for this endpoint's container. * @return the schema type. */ - SchemaType getSchemaType(); + default SchemaType getSchemaType() { + return null; + } /** * Return the concurrency for this endpoint's container. * @return the concurrency. */ @Nullable - Integer getConcurrency(); + default Integer getConcurrency() { + return null; + } /** * Setup the specified message listener container with the model defined by this @@ -101,6 +119,7 @@ public interface ListenerEndpoint { * @param listenerContainer the listener container to configure * @param messageConverter the message converter - can be null */ - void setupListenerContainer(C listenerContainer, @Nullable MessageConverter messageConverter); + default void setupListenerContainer(C listenerContainer, @Nullable MessageConverter messageConverter) { + } } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java index 2accac00..b07613b0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java @@ -28,13 +28,20 @@ import org.springframework.pulsar.listener.PulsarMessageListenerContainer; * * @author Soby Chacko * @author Alexander Preuß + * @author Vedran Pavic */ public interface PulsarListenerEndpoint extends ListenerEndpoint { - boolean isBatchListener(); + default boolean isBatchListener() { + return false; + } - Properties getConsumerProperties(); + default Properties getConsumerProperties() { + return null; + } - AckMode getAckMode(); + default AckMode getAckMode() { + return null; + } } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java index 77e9aa07..38b04a08 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java @@ -33,7 +33,9 @@ import org.springframework.pulsar.support.MessageConverter; * * @author Soby Chacko * @author Alexander Preuß + * @deprecated for removal in favor of {@link PulsarListenerEndpoint} */ +@Deprecated(forRemoval = true) public class PulsarListenerEndpointAdapter implements PulsarListenerEndpoint { @Override