public abstract class BaseHotStreamImpl<T> extends IteratorHotStream<T> implements HotStream<T>
| Modifier and Type | Field and Description |
|---|---|
protected java.util.stream.Stream<T> |
stream |
connected, connections, open, pause| Constructor and Description |
|---|
BaseHotStreamImpl(java.util.stream.Stream<T> stream) |
| Modifier and Type | Method and Description |
|---|---|
ReactiveSeq<T> |
connect(java.util.Queue<T> queue)
Connect to this HotStream using the provided transfer async.Queue.
|
abstract HotStream<T> |
init(java.util.concurrent.Executor exec) |
HotStream<T> |
paused(java.util.concurrent.Executor exec) |
HotStream<T> |
schedule(java.lang.String cron,
java.util.concurrent.ScheduledExecutorService ex) |
HotStream<T> |
scheduleFixedDelay(long delay,
java.util.concurrent.ScheduledExecutorService ex) |
HotStream<T> |
scheduleFixedRate(long rate,
java.util.concurrent.ScheduledExecutorService ex) |
isPaused, pause, scheduleFixedDelayInternal, scheduleFixedRate, scheduleInternal, unpauseprotected final java.util.stream.Stream<T> stream
public BaseHotStreamImpl(java.util.stream.Stream<T> stream)
public HotStream<T> schedule(java.lang.String cron, java.util.concurrent.ScheduledExecutorService ex)
public HotStream<T> scheduleFixedDelay(long delay, java.util.concurrent.ScheduledExecutorService ex)
public HotStream<T> scheduleFixedRate(long rate, java.util.concurrent.ScheduledExecutorService ex)
public ReactiveSeq<T> connect(java.util.Queue<T> queue)
HotStreamWaitStrategy
new LazyReact().range(0,Integer.MAX_VALUE)
.limit(1000)
.peek(v->value=v)
.peek(v->latch.countDown())
.hotStream(exec)
.connect(new LinkedBlockingQueue<>())
.limit(100)
.futureOperations(ForkJoinPool.commonPool())
.forEach(System.out::println)