DATACMNS-836 - Move reactive wrapper converters to ReactiveWrapperConverters.

Move reactive wrapper conversion to ReactiveWrapperConverters to minimize touching points with reactive APIs. ReactiveWrapperConverters is used only from reactive repository components and does not initialize with blocking API use. This removes the need of having Project Reactor on the class path for non-reactive use.
This commit is contained in:
Mark Paluch
2016-11-08 10:48:58 +01:00
committed by Oliver Gierke
parent dd2fe0a34d
commit a6bc800138
5 changed files with 672 additions and 661 deletions

View File

@@ -23,6 +23,7 @@ import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.core.convert.support.GenericConversionService;
import org.springframework.data.repository.util.NullableWrapper;
import org.springframework.data.repository.util.QueryExecutionConverters;
import org.springframework.data.repository.util.ReactiveWrapperConverters;
/**
* Simple domain service to convert query results into a dedicated type.
@@ -43,6 +44,7 @@ class QueryExecutionResultHandler {
GenericConversionService conversionService = new DefaultConversionService();
QueryExecutionConverters.registerConvertersIn(conversionService);
ReactiveWrapperConverters.registerConvertersIn(conversionService);
this.conversionService = conversionService;
}

View File

@@ -21,14 +21,11 @@ import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Future;
import org.reactivestreams.Publisher;
import org.springframework.core.ReactiveAdapterRegistry;
import org.springframework.core.convert.ConversionService;
import org.springframework.core.convert.TypeDescriptor;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.convert.converter.GenericConverter;
import org.springframework.core.convert.support.ConfigurableConversionService;
import org.springframework.data.repository.util.ReactiveWrappers.ReactiveLibrary;
import org.springframework.scheduling.annotation.AsyncResult;
import org.springframework.util.Assert;
import org.springframework.util.ClassUtils;
@@ -36,14 +33,6 @@ import org.springframework.util.concurrent.ListenableFuture;
import com.google.common.base.Optional;
import io.reactivex.BackpressureStrategy;
import io.reactivex.Flowable;
import io.reactivex.Maybe;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import rx.Completable;
import rx.Observable;
import rx.Single;
import scala.Function0;
import scala.Option;
import scala.runtime.AbstractFunction0;
@@ -85,7 +74,6 @@ public abstract class QueryExecutionConverters {
private static final Set<Class<?>> WRAPPER_TYPES = new HashSet<Class<?>>();
private static final Set<Class<?>> UNWRAPPER_TYPES = new HashSet<Class<?>>();
private static final Set<Converter<Object, Object>> UNWRAPPERS = new HashSet<Converter<Object, Object>>();
private static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry();
static {
@@ -117,9 +105,11 @@ public abstract class QueryExecutionConverters {
UNWRAPPERS.add(ScalOptionUnwrapper.INSTANCE);
}
WRAPPER_TYPES.addAll(ReactiveWrappers.getNoValueTypes());
WRAPPER_TYPES.addAll(ReactiveWrappers.getSingleValueTypes());
WRAPPER_TYPES.addAll(ReactiveWrappers.getMultiValueTypes());
if(ReactiveWrappers.isAvailable()) {
WRAPPER_TYPES.addAll(ReactiveWrappers.getNoValueTypes());
WRAPPER_TYPES.addAll(ReactiveWrappers.getSingleValueTypes());
WRAPPER_TYPES.addAll(ReactiveWrappers.getMultiValueTypes());
}
}
private QueryExecutionConverters() {}
@@ -187,67 +177,6 @@ public abstract class QueryExecutionConverters {
if (ASYNC_RESULT_PRESENT) {
conversionService.addConverter(new NullableWrapperToFutureConverter(conversionService));
}
if (ReactiveWrappers.isAvailable()) {
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) {
conversionService.addConverter(PublisherToRxJava1CompletableConverter.INSTANCE);
conversionService.addConverter(RxJava1CompletableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1CompletableToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava1SingleConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava1ObservableConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToFluxConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) {
conversionService.addConverter(PublisherToRxJava2CompletableConverter.INSTANCE);
conversionService.addConverter(RxJava2CompletableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2CompletableToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2SingleConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2ObservableConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2FlowableConverter.INSTANCE);
conversionService.addConverter(RxJava2FlowableToPublisherConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2MaybeConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToFluxConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.PROJECT_REACTOR)) {
conversionService.addConverter(PublisherToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToFluxConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) {
conversionService.addConverter(RxJava1SingleToObservableConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToSingleConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) {
conversionService.addConverter(RxJava2SingleToObservableConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToSingleConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToMaybeConverter.INSTANCE);
}
}
}
/**
@@ -557,570 +486,4 @@ public abstract class QueryExecutionConverters {
return source instanceof Option ? ((Option<?>) source).getOrElse(alternative) : source;
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToFluxConverter implements Converter<Publisher<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Publisher<?> source) {
return Flux.from(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToMonoConverter implements Converter<Publisher<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Publisher<?> source) {
return Mono.from(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1SingleConverter implements Converter<Publisher<?>, Single<?>> {
INSTANCE;
@Override
public Single<?> convert(Publisher<?> source) {
return (Single<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Single.class).fromPublisher(Mono.from(source));
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Completable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1CompletableConverter implements Converter<Publisher<?>, Completable> {
INSTANCE;
@Override
public Completable convert(Publisher<?> source) {
return (Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Completable.class).fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Observable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1ObservableConverter implements Converter<Publisher<?>, Observable<?>> {
INSTANCE;
@Override
public Observable<?> convert(Publisher<?> source) {
return (Observable<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Observable.class).fromPublisher(Flux.from(source));
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToPublisherConverter implements Converter<Single<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToMonoConverter implements Converter<Single<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToFluxConverter implements Converter<Single<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Completable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1CompletableToPublisherConverter implements Converter<Completable, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Completable source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Completable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Completable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1CompletableToMonoConverter implements Converter<Completable, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Completable source) {
return Mono.from(RxJava1CompletableToPublisherConverter.INSTANCE.convert(source));
}
}
/**
* A {@link Converter} to convert an {@link Observable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToPublisherConverter implements Converter<Observable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToMonoConverter implements Converter<Observable<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToFluxConverter implements Converter<Observable<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2SingleConverter implements Converter<Publisher<?>, io.reactivex.Single<?>> {
INSTANCE;
@Override
public io.reactivex.Single<?> convert(Publisher<?> source) {
return (io.reactivex.Single<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Completable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2CompletableConverter implements Converter<Publisher<?>, io.reactivex.Completable> {
INSTANCE;
@Override
public io.reactivex.Completable convert(Publisher<?> source) {
return (io.reactivex.Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Completable.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Observable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2ObservableConverter implements Converter<Publisher<?>, io.reactivex.Observable<?>> {
INSTANCE;
@Override
public io.reactivex.Observable<?> convert(Publisher<?> source) {
return (io.reactivex.Observable<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToPublisherConverter implements Converter<io.reactivex.Single<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToMonoConverter implements Converter<io.reactivex.Single<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToFluxConverter implements Converter<io.reactivex.Single<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2CompletableToPublisherConverter implements Converter<io.reactivex.Completable, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Completable source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Completable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2CompletableToMonoConverter implements Converter<io.reactivex.Completable, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Completable source) {
return Mono.from(RxJava2CompletableToPublisherConverter.INSTANCE.convert(source));
}
}
/**
* A {@link Converter} to convert an {@link io.reactivex.Observable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToPublisherConverter implements Converter<io.reactivex.Observable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Observable<?> source) {
return source.toFlowable(BackpressureStrategy.BUFFER);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToMonoConverter implements Converter<io.reactivex.Observable<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Observable<?> source) {
return Mono.from(source.toFlowable(BackpressureStrategy.BUFFER));
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToFluxConverter implements Converter<io.reactivex.Observable<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Observable<?> source) {
return Flux.from(source.toFlowable(BackpressureStrategy.BUFFER));
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2FlowableConverter implements Converter<Publisher<?>, io.reactivex.Flowable<?>> {
INSTANCE;
@Override
public io.reactivex.Flowable<?> convert(Publisher<?> source) {
return Flowable.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Flowable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2FlowableToPublisherConverter implements Converter<io.reactivex.Flowable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Flowable<?> source) {
return source;
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2MaybeConverter implements Converter<Publisher<?>, io.reactivex.Maybe<?>> {
INSTANCE;
@Override
public io.reactivex.Maybe<?> convert(Publisher<?> source) {
return (io.reactivex.Maybe<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Maybe.class).fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToPublisherConverter implements Converter<io.reactivex.Maybe<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Maybe<?> source) {
return source.toFlowable();
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToMonoConverter implements Converter<io.reactivex.Maybe<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Maybe<?> source) {
return Mono.from(source.toFlowable());
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToFluxConverter implements Converter<io.reactivex.Maybe<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Maybe<?> source) {
return Flux.from(source.toFlowable());
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToSingleConverter implements Converter<Observable<?>, Single<?>> {
INSTANCE;
@Override
public Single<?> convert(Observable<?> source) {
return source.toSingle();
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToObservableConverter implements Converter<Single<?>, Observable<?>> {
INSTANCE;
@Override
public Observable<?> convert(Single<?> source) {
return source.toObservable();
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToSingleConverter
implements Converter<io.reactivex.Observable<?>, io.reactivex.Single<?>> {
INSTANCE;
@Override
public io.reactivex.Single<?> convert(io.reactivex.Observable<?> source) {
return source.singleOrError();
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Maybe}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToMaybeConverter
implements Converter<io.reactivex.Observable<?>, io.reactivex.Maybe<?>> {
INSTANCE;
@Override
public io.reactivex.Maybe<?> convert(io.reactivex.Observable<?> source) {
return source.singleElement();
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToObservableConverter
implements Converter<io.reactivex.Single<?>, io.reactivex.Observable<?>> {
INSTANCE;
@Override
public io.reactivex.Observable<?> convert(io.reactivex.Single<?> source) {
return source.toObservable();
}
}
}

View File

@@ -15,26 +15,32 @@
*/
package org.springframework.data.repository.util;
import static org.springframework.data.repository.util.ReactiveWrapperConverters.RegistryHolder.*;
import java.util.ArrayList;
import java.util.List;
import org.reactivestreams.Publisher;
import org.springframework.core.ReactiveAdapterRegistry;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.convert.support.ConfigurableConversionService;
import org.springframework.core.convert.support.GenericConversionService;
import org.springframework.data.repository.util.ReactiveWrappers.ReactiveLibrary;
import org.springframework.util.Assert;
import org.springframework.util.ClassUtils;
import io.reactivex.BackpressureStrategy;
import io.reactivex.Flowable;
import io.reactivex.Maybe;
import lombok.experimental.UtilityClass;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import rx.Completable;
import rx.Observable;
import rx.Single;
/**
* Conversion support for reactive wrapper types. This class is a logical extension to {@link QueryExecutionConverters}.
* Conversion support for reactive wrapper types. This class is a reactive extension to {@link QueryExecutionConverters}.
* <p>
* This class discovers reactive wrapper availability and their conversion support based on the class path. Reactive
* wrapper types might be supported/on the class path but conversion may require additional dependencies.
@@ -49,7 +55,6 @@ public class ReactiveWrapperConverters {
private static final List<AbstractReactiveWrapper<?>> REACTIVE_WRAPPERS = new ArrayList<>();
private static final GenericConversionService GENERIC_CONVERSION_SERVICE = new GenericConversionService();
private static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry();
static {
@@ -74,7 +79,76 @@ public class ReactiveWrapperConverters {
REACTIVE_WRAPPERS.add(PublisherWrapper.INSTANCE);
}
QueryExecutionConverters.registerConvertersIn(GENERIC_CONVERSION_SERVICE);
registerConvertersIn(GENERIC_CONVERSION_SERVICE);
}
/**
* Registers converters for wrapper types found on the classpath.
*
* @param conversionService must not be {@literal null}.
*/
public static void registerConvertersIn(ConfigurableConversionService conversionService) {
Assert.notNull(conversionService, "ConversionService must not be null!");
if (ReactiveWrappers.isAvailable(ReactiveLibrary.PROJECT_REACTOR)) {
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) {
conversionService.addConverter(PublisherToRxJava1CompletableConverter.INSTANCE);
conversionService.addConverter(RxJava1CompletableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1CompletableToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava1SingleConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava1SingleToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava1ObservableConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToFluxConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) {
conversionService.addConverter(PublisherToRxJava2CompletableConverter.INSTANCE);
conversionService.addConverter(RxJava2CompletableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2CompletableToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2SingleConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2SingleToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2ObservableConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToFluxConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2FlowableConverter.INSTANCE);
conversionService.addConverter(RxJava2FlowableToPublisherConverter.INSTANCE);
conversionService.addConverter(PublisherToRxJava2MaybeConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToPublisherConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToMonoConverter.INSTANCE);
conversionService.addConverter(RxJava2MaybeToFluxConverter.INSTANCE);
}
conversionService.addConverter(PublisherToMonoConverter.INSTANCE);
conversionService.addConverter(PublisherToFluxConverter.INSTANCE);
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) {
conversionService.addConverter(RxJava1SingleToObservableConverter.INSTANCE);
conversionService.addConverter(RxJava1ObservableToSingleConverter.INSTANCE);
}
if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) {
conversionService.addConverter(RxJava2SingleToObservableConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToSingleConverter.INSTANCE);
conversionService.addConverter(RxJava2ObservableToMaybeConverter.INSTANCE);
}
}
}
/**
@@ -88,7 +162,7 @@ public class ReactiveWrapperConverters {
* @return {@literal true} if the {@code type} is a supported reactive wrapper type.
*/
public static boolean supports(Class<?> type) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(type) != null;
return RegistryHolder.REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(type) != null;
}
/**
@@ -283,4 +357,580 @@ public class ReactiveWrapperConverters {
return ((io.reactivex.Flowable<?>) wrapper).map(converter::convert);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToFluxConverter implements Converter<Publisher<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Publisher<?> source) {
return Flux.from(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToMonoConverter implements Converter<Publisher<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Publisher<?> source) {
return Mono.from(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1SingleConverter implements Converter<Publisher<?>, Single<?>> {
INSTANCE;
@Override
public Single<?> convert(Publisher<?> source) {
return (Single<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Single.class).fromPublisher(Mono.from(source));
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Completable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1CompletableConverter implements Converter<Publisher<?>, Completable> {
INSTANCE;
@Override
public Completable convert(Publisher<?> source) {
return (Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Completable.class).fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link Observable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava1ObservableConverter implements Converter<Publisher<?>, Observable<?>> {
INSTANCE;
@Override
public Observable<?> convert(Publisher<?> source) {
return (Observable<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Observable.class).fromPublisher(Flux.from(source));
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToPublisherConverter implements Converter<Single<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToMonoConverter implements Converter<Single<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToFluxConverter implements Converter<Single<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Completable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1CompletableToPublisherConverter implements Converter<Completable, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Completable source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Completable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Completable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1CompletableToMonoConverter implements Converter<Completable, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Completable source) {
return Mono.from(RxJava1CompletableToPublisherConverter.INSTANCE.convert(source));
}
}
/**
* A {@link Converter} to convert an {@link Observable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToPublisherConverter implements Converter<Observable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToMonoConverter implements Converter<Observable<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToFluxConverter implements Converter<Observable<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(Observable<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2SingleConverter implements Converter<Publisher<?>, io.reactivex.Single<?>> {
INSTANCE;
@Override
public io.reactivex.Single<?> convert(Publisher<?> source) {
return (io.reactivex.Single<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Completable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2CompletableConverter implements Converter<Publisher<?>, io.reactivex.Completable> {
INSTANCE;
@Override
public io.reactivex.Completable convert(Publisher<?> source) {
return (io.reactivex.Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Completable.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Observable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2ObservableConverter implements Converter<Publisher<?>, io.reactivex.Observable<?>> {
INSTANCE;
@Override
public io.reactivex.Observable<?> convert(Publisher<?> source) {
return (io.reactivex.Observable<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class)
.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToPublisherConverter implements Converter<io.reactivex.Single<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToMonoConverter implements Converter<io.reactivex.Single<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToFluxConverter implements Converter<io.reactivex.Single<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Single<?> source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2CompletableToPublisherConverter implements Converter<io.reactivex.Completable, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Completable source) {
return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Completable.class).toFlux(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2CompletableToMonoConverter implements Converter<io.reactivex.Completable, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Completable source) {
return Mono.from(RxJava2CompletableToPublisherConverter.INSTANCE.convert(source));
}
}
/**
* A {@link Converter} to convert an {@link io.reactivex.Observable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToPublisherConverter implements Converter<io.reactivex.Observable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Observable<?> source) {
return source.toFlowable(BackpressureStrategy.BUFFER);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToMonoConverter implements Converter<io.reactivex.Observable<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Observable<?> source) {
return Mono.from(source.toFlowable(BackpressureStrategy.BUFFER));
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToFluxConverter implements Converter<io.reactivex.Observable<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Observable<?> source) {
return Flux.from(source.toFlowable(BackpressureStrategy.BUFFER));
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2FlowableConverter implements Converter<Publisher<?>, io.reactivex.Flowable<?>> {
INSTANCE;
@Override
public io.reactivex.Flowable<?> convert(Publisher<?> source) {
return Flowable.fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Flowable} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2FlowableToPublisherConverter implements Converter<io.reactivex.Flowable<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Flowable<?> source) {
return source;
}
}
/**
* A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum PublisherToRxJava2MaybeConverter implements Converter<Publisher<?>, io.reactivex.Maybe<?>> {
INSTANCE;
@Override
public io.reactivex.Maybe<?> convert(Publisher<?> source) {
return (io.reactivex.Maybe<?>) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Maybe.class).fromPublisher(source);
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Publisher}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToPublisherConverter implements Converter<io.reactivex.Maybe<?>, Publisher<?>> {
INSTANCE;
@Override
public Publisher<?> convert(io.reactivex.Maybe<?> source) {
return source.toFlowable();
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Mono}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToMonoConverter implements Converter<io.reactivex.Maybe<?>, Mono<?>> {
INSTANCE;
@Override
public Mono<?> convert(io.reactivex.Maybe<?> source) {
return Mono.from(source.toFlowable());
}
}
/**
* A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Flux}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2MaybeToFluxConverter implements Converter<io.reactivex.Maybe<?>, Flux<?>> {
INSTANCE;
@Override
public Flux<?> convert(io.reactivex.Maybe<?> source) {
return Flux.from(source.toFlowable());
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1ObservableToSingleConverter implements Converter<Observable<?>, Single<?>> {
INSTANCE;
@Override
public Single<?> convert(Observable<?> source) {
return source.toSingle();
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava1SingleToObservableConverter implements Converter<Single<?>, Observable<?>> {
INSTANCE;
@Override
public Observable<?> convert(Single<?> source) {
return source.toObservable();
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToSingleConverter
implements Converter<io.reactivex.Observable<?>, io.reactivex.Single<?>> {
INSTANCE;
@Override
public io.reactivex.Single<?> convert(io.reactivex.Observable<?> source) {
return source.singleOrError();
}
}
/**
* A {@link Converter} to convert a {@link Observable} to {@link Maybe}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2ObservableToMaybeConverter
implements Converter<io.reactivex.Observable<?>, io.reactivex.Maybe<?>> {
INSTANCE;
@Override
public io.reactivex.Maybe<?> convert(io.reactivex.Observable<?> source) {
return source.singleElement();
}
}
/**
* A {@link Converter} to convert a {@link Single} to {@link Single}.
*
* @author Mark Paluch
* @author 2.0
*/
public enum RxJava2SingleToObservableConverter
implements Converter<io.reactivex.Single<?>, io.reactivex.Observable<?>> {
INSTANCE;
@Override
public io.reactivex.Observable<?> convert(io.reactivex.Single<?> source) {
return source.toObservable();
}
}
/**
* Holder for delayed initialization of {@link ReactiveAdapterRegistry}.
*
* @author Mark Paluch
* @author 2.0
*/
static class RegistryHolder {
static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry();
}
}