public class FluxSource
extends java.lang.Object
| Modifier and Type | Method and Description |
|---|---|
<T> PushableFlux<T> |
flux()
Create a pushable Flux
|
static <T> reactor.core.publisher.Flux<T> |
flux(Adapter<T> adapter)
Create a pushable ReactiveSeq
|
static <T> LazyFutureStream<T> |
futureStream(Adapter<T> adapter,
LazyReact react)
Create a LazyFutureStream.
|
<T> PushableLazyFutureStream<T> |
futureStream(LazyReact s)
Create a pushable LazyFutureStream using the supplied ReactPool
|
static FluxSource |
of(int backPressureAfter)
Create a Pushable Flux soure backed by a bounded BlockingQueue.
|
static FluxSource |
of(QueueFactory<?> q)
Create a Pushable Flux source backed by a queue created by the supplied queue factory
|
static <T> MultipleFluxSource<T> |
ofMultiple() |
static <T> MultipleFluxSource<T> |
ofMultiple(int backPressureAfter) |
static <T> MultipleFluxSource<T> |
ofMultiple(QueueFactory<?> q) |
static FluxSource |
ofUnbounded()
Create a Pushable Flux source backed by an unbounded ConcurrentLinkedQueue
|
<T> PushableReactiveSeq<T> |
reactiveSeq()
Create a pushable ReactiveSeq
|
static <T> ReactiveSeq<T> |
reactiveSeq(Adapter<T> adapter)
Create a pushable ReactiveSeq
|
<T> PushableStream<T> |
stream()
Create a pushable JDK 8 Stream
|
static <T> java.util.stream.Stream<T> |
stream(Adapter<T> adapter)
Create a JDK 8 Stream from the supplied Adapter
|
public static <T> MultipleFluxSource<T> ofMultiple()
public static <T> MultipleFluxSource<T> ofMultiple(int backPressureAfter)
public static <T> MultipleFluxSource<T> ofMultiple(QueueFactory<?> q)
public static FluxSource of(QueueFactory<?> q)
//create a QueueFactory for Agrona OneToOneConcurrentArrayQueue with capacity for 1,000 elements
QueueFactory<Integer> input = QueueFactories.singleWriterboundedNonBlockingQueue(1_000);
FluxSource source = FluxSource.of(input);
PushableFlux flux = source.flux();
flux.getInput().offer(1);
//on a different thread
flux.getFlux()
.map(i->i*2)
.subscribe(System.out::println);
q - Queue Factory used to provide data for the Flux streampublic static FluxSource ofUnbounded()
public static FluxSource of(int backPressureAfter)
backPressureAfter - Max queue size - once queue is full, producers adding to the Queue will be blockedpublic <T> PushableLazyFutureStream<T> futureStream(LazyReact s)
s - ReactPool to use to create the Streampublic static <T> LazyFutureStream<T> futureStream(Adapter<T> adapter, LazyReact react)
adapter - Adapter to create a LazyFutureStream frompublic <T> PushableStream<T> stream()
public <T> PushableFlux<T> flux()
public <T> PushableReactiveSeq<T> reactiveSeq()
public static <T> java.util.stream.Stream<T> stream(Adapter<T> adapter)
adapter - Adapter to create a Steam frompublic static <T> ReactiveSeq<T> reactiveSeq(Adapter<T> adapter)
adapter - Adapter to create a Seq frompublic static <T> reactor.core.publisher.Flux<T> flux(Adapter<T> adapter)
adapter - Adapter to create a Flux from