add interval for non-Flux stream Suppliers

This commit is contained in:
markfisher
2017-02-24 12:44:14 -05:00
parent 1dc0489c23
commit 2ae7789cd1
5 changed files with 38 additions and 5 deletions

View File

@@ -18,6 +18,7 @@ package org.springframework.cloud.function.support;
import java.time.Duration;
import java.util.function.Supplier;
import java.util.stream.Stream;
import reactor.core.publisher.Flux;
@@ -48,10 +49,18 @@ public class FluxSupplier<T> implements Supplier<Flux<T>> {
}
@Override
@SuppressWarnings({ "unchecked", "rawtypes" })
public Flux<T> get() {
if (this.period != null) {
return Flux.interval(this.period).map(i->this.supplier.get());
}
return Flux.just(this.supplier.get());
Object result = this.supplier.get();
if (result instanceof Stream) {
return Flux.fromStream((Stream) result);
}
if (result instanceof Iterable) {
return Flux.fromIterable((Iterable) result);
}
return Flux.just((T) result);
}
}