-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Add a general purpose collecting method to ReadStream to facilitate t…
…he reduction of streams. The collecting method is a default method.
- Loading branch information
Showing
4 changed files
with
110 additions
and
2 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -11,7 +11,7 @@ | |
|
||
package examples; | ||
|
||
import io.vertx.core.Handler; | ||
import io.vertx.core.Future; | ||
import io.vertx.core.Vertx; | ||
import io.vertx.core.buffer.Buffer; | ||
import io.vertx.core.file.AsyncFile; | ||
|
@@ -20,6 +20,9 @@ | |
import io.vertx.core.net.NetServer; | ||
import io.vertx.core.net.NetServerOptions; | ||
import io.vertx.core.streams.Pipe; | ||
import io.vertx.core.streams.ReadStream; | ||
|
||
import java.util.stream.Collectors; | ||
|
||
/** | ||
* @author <a href="mailto:[email protected]">Julien Viet</a> | ||
|
@@ -164,4 +167,11 @@ public void pipe9(AsyncFile src, AsyncFile dst) { | |
dst.end(Buffer.buffer("done")); | ||
}); | ||
} | ||
|
||
public <T> void reduce1(ReadStream<T> stream) { | ||
// Count the number of elements | ||
Future<Long> result = stream.collect(Collectors.counting()); | ||
|
||
result.onSuccess(count -> System.out.println("Stream emitted " + count + " elements")); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
62 changes: 62 additions & 0 deletions
62
src/test/java/io/vertx/core/streams/ReadStreamReduceTest.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,62 @@ | ||
/* | ||
* Copyright (c) 2011-2019 Contributors to the Eclipse Foundation | ||
* | ||
* This program and the accompanying materials are made available under the | ||
* terms of the Eclipse Public License 2.0 which is available at | ||
* http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 | ||
* which is available at https://www.apache.org/licenses/LICENSE-2.0. | ||
* | ||
* SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 | ||
*/ | ||
package io.vertx.core.streams; | ||
|
||
import io.vertx.core.Future; | ||
import io.vertx.test.core.AsyncTestBase; | ||
import io.vertx.test.fakestream.FakeStream; | ||
import org.junit.Test; | ||
|
||
import java.util.Arrays; | ||
import java.util.List; | ||
import java.util.stream.Collectors; | ||
|
||
public class ReadStreamReduceTest extends AsyncTestBase { | ||
|
||
private FakeStream<Object> dst; | ||
private Object o1 = new Object(); | ||
private Object o2 = new Object(); | ||
private Object o3 = new Object(); | ||
|
||
@Override | ||
protected void setUp() throws Exception { | ||
super.setUp(); | ||
dst = new FakeStream<>(); | ||
} | ||
|
||
@Test | ||
public void testCollect() { | ||
Future<List<Object>> list = dst.collect(Collectors.toList()); | ||
assertFalse(list.isComplete()); | ||
dst.write(o1); | ||
assertFalse(list.isComplete()); | ||
dst.write(o2); | ||
assertFalse(list.isComplete()); | ||
dst.write(o3); | ||
dst.end(); | ||
assertTrue(list.succeeded()); | ||
assertEquals(Arrays.asList(o1, o2, o3), list.result()); | ||
} | ||
|
||
@Test | ||
public void testFailure() { | ||
Future<List<Object>> list = dst.collect(Collectors.toList()); | ||
assertFalse(list.isComplete()); | ||
dst.write(o1); | ||
assertFalse(list.isComplete()); | ||
dst.write(o2); | ||
assertFalse(list.isComplete()); | ||
Throwable err = new Throwable(); | ||
dst.fail(err); | ||
assertTrue(list.failed()); | ||
assertSame(err, list.cause()); | ||
} | ||
} |