[英]Return a standard (unextended) JDK Stream connected to this Queue To disconnect cleanly close the queue
use queue.stream().parallel() to convert to a parallel Stream
use queue.stream().parallel() to convert to a parallel Stream
代码示例来源:origin: aol/cyclops
* Return a standard (unextended) JDK Stream connected to this Queue
* To disconnect cleanly close the queue
* <pre>
* {@code
* use queue.stream().parallel() to convert to a parallel Stream
* }
* </pre>
* @see Queue#jdkStream(int) for an alternative that sends more poision pills for use with parallel Streams.
* @return Java 8 Stream connnected to this Queue
public Stream<T> jdkStream() {
return jdkStream(2);
代码示例来源:origin: aol/cyclops
public Iterator<T> iterator() {
return host.jdkStream().iterator();
代码示例来源:origin: aol/cyclops
default <R> R foldParallel(Function<? super Stream<T>,? extends R> fn){
Queue<T> queue = QueueFactories.<T>unboundedNonBlockingQueue().build().withTimeout(1);
AtomicReference<Continuation> ref = new AtomicReference<>(null);
Continuation cont =
new Continuation(()->{
if(ref.get()==null && ref.compareAndSet(null,Continuation.empty())){
try {
//use the first consuming thread to tell this Stream onto the Queue
}finally {
return Continuation.empty();
return fn.apply(queue.jdkStream().parallel());
default <R> R foldParallel(ForkJoinPool fj,Function<? super Stream<T>,? extends R> fn){
代码示例来源:origin: com.oath.cyclops/cyclops
* Return a standard (unextended) JDK Stream connected to this Queue
* To disconnect cleanly close the queue
* <pre>
* {@code
* use queue.stream().parallel() to convert to a parallel Stream
* }
* </pre>
* @see Queue#jdkStream(int) for an alternative that sends more poision pills for use with parallel Streams.
* @return Java 8 Stream connnected to this Queue
public Stream<T> jdkStream() {
return jdkStream(2);
代码示例来源:origin: com.oath.cyclops/cyclops
public Iterator<T> iterator() {
return host.jdkStream().iterator();
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> forkInput = queue.jdkStream();
for (OperationProcessor opProcInFork : opProcsInFork) {
forkInput = opProcInFork.perform(forkInput);
return outputQueue.jdkStream();
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> forkInput = queue.jdkStream();
for (OperationProcessor opProcInFork : opProcsInFork) {
forkInput = opProcInFork.perform(forkInput);
return outputQueue.jdkStream();
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> conditionInput = queue.jdkStream();
for (OperationProcessor proc : procs) {
conditionInput = proc.perform(conditionInput);
return outputQueue.jdkStream();
代码示例来源:origin: com.oath.cyclops/cyclops-futurestream
* Create a pushable JDK 8 Stream
* <pre>
* {@code
* PushableStream<Integer> pushable = StreamSource.ofUnbounded()
pushable.getStream().collect(CyclopsCollectors.toList()) //[10]
* }
* </pre>
* @return PushableStream that can accept data to push into a Java 8 Stream
* to push it to the Stream
public <T> PushableStream<T> stream() {
final Queue<T> q = createQueue();
return new PushableStream<T>(
q, q.jdkStream());
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> conditionInput = queue.jdkStream();
for (OperationProcessor proc : procs) {
conditionInput = proc.perform(conditionInput);
return outputQueue.jdkStream();
代码示例来源:origin: com.oath.cyclops/cyclops-futurestream
default ReactiveSeq<U> stream() {
return Streams.oneShotStream(toQueue().jdkStream(getSubscription()));
代码示例来源:origin: com.oath.cyclops/cyclops
default <R> R foldParallel(Function<? super Stream<T>,? extends R> fn){
Queue<T> queue = QueueFactories.<T>unboundedNonBlockingQueue().build().withTimeout(1);
AtomicReference<Continuation> ref = new AtomicReference<>(null);
Continuation cont =
new Continuation(()->{
if(ref.get()==null && ref.compareAndSet(null,Continuation.empty())){
try {
//use the first consuming thread to tell this Stream onto the Queue
}finally {
return Continuation.empty();
return fn.apply(queue.jdkStream().parallel());
default <R> R foldParallel(ForkJoinPool fj,Function<? super Stream<T>,? extends R> fn){
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> input = this.eventQueue.jdkStream();
代码示例来源:origin: Nextdoor/bender
Stream<InternalEvent> input = this.eventQueue.jdkStream();