[BEAM-443] Update Beam batch examples to call waitUntilFinish()
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/e5afbe56 Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/e5afbe56 Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/e5afbe56 Branch: refs/heads/apex-runner Commit: e5afbe560b604ae7081e420af5b0d855508d53ad Parents: eba099f Author: Pei He <pe...@google.com> Authored: Thu Oct 13 14:44:13 2016 -0700 Committer: Pei He <pe...@google.com> Committed: Wed Oct 26 16:02:17 2016 -0700 ---------------------------------------------------------------------- .../java/org/apache/beam/examples/DebuggingWordCount.java | 2 +- .../java/org/apache/beam/examples/MinimalWordCount.java | 2 +- .../src/main/java/org/apache/beam/examples/WordCount.java | 2 +- .../java/org/apache/beam/examples/complete/TfIdf.java | 2 +- .../beam/examples/complete/TopWikipediaSessions.java | 2 +- .../apache/beam/examples/cookbook/BigQueryTornadoes.java | 2 +- .../beam/examples/cookbook/CombinePerKeyExamples.java | 2 +- .../org/apache/beam/examples/cookbook/DeDupExample.java | 2 +- .../org/apache/beam/examples/cookbook/FilterExamples.java | 2 +- .../org/apache/beam/examples/cookbook/JoinExamples.java | 2 +- .../apache/beam/examples/cookbook/MaxPerKeyExamples.java | 2 +- .../test/java/org/apache/beam/examples/WordCountTest.java | 2 +- .../apache/beam/examples/complete/AutoCompleteTest.java | 6 +++--- .../java/org/apache/beam/examples/complete/TfIdfTest.java | 2 +- .../beam/examples/complete/TopWikipediaSessionsTest.java | 2 +- .../apache/beam/examples/cookbook/DeDupExampleTest.java | 4 ++-- .../apache/beam/examples/cookbook/JoinExamplesTest.java | 2 +- .../apache/beam/examples/cookbook/TriggerExampleTest.java | 2 +- .../org/apache/beam/examples/MinimalWordCountJava8.java | 2 +- .../beam/examples/complete/game/HourlyTeamScore.java | 2 +- .../org/apache/beam/examples/complete/game/UserScore.java | 2 +- .../apache/beam/examples/complete/game/GameStatsTest.java | 2 +- .../beam/examples/complete/game/HourlyTeamScoreTest.java | 2 +- .../beam/examples/complete/game/LeaderBoardTest.java | 10 +++++----- .../apache/beam/examples/complete/game/UserScoreTest.java | 6 +++--- 25 files changed, 34 insertions(+), 34 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/DebuggingWordCount.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/DebuggingWordCount.java b/examples/java/src/main/java/org/apache/beam/examples/DebuggingWordCount.java index 90d77b3..1d2c83a 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/DebuggingWordCount.java +++ b/examples/java/src/main/java/org/apache/beam/examples/DebuggingWordCount.java @@ -194,6 +194,6 @@ public class DebuggingWordCount { KV.of("stomach", 1L)); PAssert.that(filteredWords).containsInAnyOrder(expectedResults); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java b/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java index 14ffa18..6fc873e 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java +++ b/examples/java/src/main/java/org/apache/beam/examples/MinimalWordCount.java @@ -119,6 +119,6 @@ public class MinimalWordCount { .apply(TextIO.Write.to("gs://YOUR_OUTPUT_BUCKET/AND_OUTPUT_PREFIX")); // Run the pipeline. - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/WordCount.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/WordCount.java b/examples/java/src/main/java/org/apache/beam/examples/WordCount.java index 1e03b34..e7eab6e 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/WordCount.java +++ b/examples/java/src/main/java/org/apache/beam/examples/WordCount.java @@ -197,6 +197,6 @@ public class WordCount { .apply(MapElements.via(new FormatAsTextFn())) .apply("WriteCounts", TextIO.Write.to(options.getOutput())); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/complete/TfIdf.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/complete/TfIdf.java b/examples/java/src/main/java/org/apache/beam/examples/complete/TfIdf.java index d4107c9..c0ba1e9 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/complete/TfIdf.java +++ b/examples/java/src/main/java/org/apache/beam/examples/complete/TfIdf.java @@ -417,6 +417,6 @@ public class TfIdf { .apply(new ComputeTfIdf()) .apply(new WriteTfIdf(options.getOutput())); - pipeline.run(); + pipeline.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/complete/TopWikipediaSessions.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/complete/TopWikipediaSessions.java b/examples/java/src/main/java/org/apache/beam/examples/complete/TopWikipediaSessions.java index 15923eb..d57cc3a 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/complete/TopWikipediaSessions.java +++ b/examples/java/src/main/java/org/apache/beam/examples/complete/TopWikipediaSessions.java @@ -208,6 +208,6 @@ public class TopWikipediaSessions { .apply(new ComputeTopSessions(samplingThreshold)) .apply("Write", TextIO.Write.withoutSharding().to(options.getOutput())); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/BigQueryTornadoes.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/BigQueryTornadoes.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/BigQueryTornadoes.java index 391ea90..a4c1a6b 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/BigQueryTornadoes.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/BigQueryTornadoes.java @@ -164,6 +164,6 @@ public class BigQueryTornadoes { .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE)); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/CombinePerKeyExamples.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/CombinePerKeyExamples.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/CombinePerKeyExamples.java index 1f0abce..93eee15 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/CombinePerKeyExamples.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/CombinePerKeyExamples.java @@ -208,6 +208,6 @@ public class CombinePerKeyExamples { .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE)); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/DeDupExample.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/DeDupExample.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/DeDupExample.java index 92f5b93..0883815 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/DeDupExample.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/DeDupExample.java @@ -91,6 +91,6 @@ public class DeDupExample { .apply(RemoveDuplicates.<String>create()) .apply("DedupedShakespeare", TextIO.Write.to(options.getOutput())); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/FilterExamples.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/FilterExamples.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/FilterExamples.java index 0b2ae73..6e6452c 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/FilterExamples.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/FilterExamples.java @@ -247,6 +247,6 @@ public class FilterExamples { .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE)); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/JoinExamples.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/JoinExamples.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/JoinExamples.java index d66e070..7cf0942 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/JoinExamples.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/JoinExamples.java @@ -170,7 +170,7 @@ public class JoinExamples { PCollection<TableRow> countryCodes = p.apply(BigQueryIO.Read.from(COUNTRY_CODES)); PCollection<String> formattedResults = joinEvents(eventsTable, countryCodes); formattedResults.apply(TextIO.Write.to(options.getOutput())); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/main/java/org/apache/beam/examples/cookbook/MaxPerKeyExamples.java ---------------------------------------------------------------------- diff --git a/examples/java/src/main/java/org/apache/beam/examples/cookbook/MaxPerKeyExamples.java b/examples/java/src/main/java/org/apache/beam/examples/cookbook/MaxPerKeyExamples.java index eed4bbd..abc10f3 100644 --- a/examples/java/src/main/java/org/apache/beam/examples/cookbook/MaxPerKeyExamples.java +++ b/examples/java/src/main/java/org/apache/beam/examples/cookbook/MaxPerKeyExamples.java @@ -157,6 +157,6 @@ public class MaxPerKeyExamples { .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED) .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE)); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/WordCountTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/WordCountTest.java b/examples/java/src/test/java/org/apache/beam/examples/WordCountTest.java index 98c5b17..c8809de 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/WordCountTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/WordCountTest.java @@ -80,6 +80,6 @@ public class WordCountTest { .apply(MapElements.via(new FormatAsTextFn())); PAssert.that(output).containsInAnyOrder(COUNTS_ARRAY); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/complete/AutoCompleteTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/complete/AutoCompleteTest.java b/examples/java/src/test/java/org/apache/beam/examples/complete/AutoCompleteTest.java index b6751c5..5dbfa70 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/complete/AutoCompleteTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/complete/AutoCompleteTest.java @@ -99,7 +99,7 @@ public class AutoCompleteTest implements Serializable { KV.of("bl", parseList("blackberry:3", "blueberry:2")), KV.of("c", parseList("cherry:1")), KV.of("ch", parseList("cherry:1"))); - p.run(); + p.run().waitUntilFinish(); } @Test @@ -117,7 +117,7 @@ public class AutoCompleteTest implements Serializable { KV.of("x", parseList("x:3", "xy:2")), KV.of("xy", parseList("xy:2", "xyz:1")), KV.of("xyz", parseList("xyz:1"))); - p.run(); + p.run().waitUntilFinish(); } @Test @@ -153,7 +153,7 @@ public class AutoCompleteTest implements Serializable { // Window [2, 3) KV.of("x", parseList("xB:2")), KV.of("xB", parseList("xB:2"))); - p.run(); + p.run().waitUntilFinish(); } private static List<CompletionCandidate> parseList(String... entries) { http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/complete/TfIdfTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/complete/TfIdfTest.java b/examples/java/src/test/java/org/apache/beam/examples/complete/TfIdfTest.java index c2d654e..1aee8f9 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/complete/TfIdfTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/complete/TfIdfTest.java @@ -61,6 +61,6 @@ public class TfIdfTest { PAssert.that(words).containsInAnyOrder(Arrays.asList("a", "m", "n", "b", "c", "d")); - pipeline.run(); + pipeline.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/complete/TopWikipediaSessionsTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/complete/TopWikipediaSessionsTest.java b/examples/java/src/test/java/org/apache/beam/examples/complete/TopWikipediaSessionsTest.java index 42fb06a..154ea73 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/complete/TopWikipediaSessionsTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/complete/TopWikipediaSessionsTest.java @@ -56,6 +56,6 @@ public class TopWikipediaSessionsTest { "user3 : [1970-02-05T00:00:00.000Z..1970-02-05T01:00:00.000Z)" + " : 1 : 1970-02-01T00:00:00.000Z")); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/cookbook/DeDupExampleTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/cookbook/DeDupExampleTest.java b/examples/java/src/test/java/org/apache/beam/examples/cookbook/DeDupExampleTest.java index c725e4f..d29fc42 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/cookbook/DeDupExampleTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/cookbook/DeDupExampleTest.java @@ -59,7 +59,7 @@ public class DeDupExampleTest { PAssert.that(output) .containsInAnyOrder("k1", "k5", "k2", "k3"); - p.run(); + p.run().waitUntilFinish(); } @Test @@ -77,6 +77,6 @@ public class DeDupExampleTest { input.apply(RemoveDuplicates.<String>create()); PAssert.that(output).empty(); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/cookbook/JoinExamplesTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/cookbook/JoinExamplesTest.java b/examples/java/src/test/java/org/apache/beam/examples/cookbook/JoinExamplesTest.java index 60f71a2..6c54aff 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/cookbook/JoinExamplesTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/cookbook/JoinExamplesTest.java @@ -108,6 +108,6 @@ public class JoinExamplesTest { PCollection<String> output = JoinExamples.joinEvents(input1, input2); PAssert.that(output).containsInAnyOrder(JOINED_EVENTS); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java/src/test/java/org/apache/beam/examples/cookbook/TriggerExampleTest.java ---------------------------------------------------------------------- diff --git a/examples/java/src/test/java/org/apache/beam/examples/cookbook/TriggerExampleTest.java b/examples/java/src/test/java/org/apache/beam/examples/cookbook/TriggerExampleTest.java index 3848ca1..bdda22c 100644 --- a/examples/java/src/test/java/org/apache/beam/examples/cookbook/TriggerExampleTest.java +++ b/examples/java/src/test/java/org/apache/beam/examples/cookbook/TriggerExampleTest.java @@ -123,7 +123,7 @@ public class TriggerExampleTest { PAssert.that(results) .containsInAnyOrder(canonicalFormat(OUT_ROW_1), canonicalFormat(OUT_ROW_2)); - pipeline.run(); + pipeline.run().waitUntilFinish(); } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/main/java/org/apache/beam/examples/MinimalWordCountJava8.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/main/java/org/apache/beam/examples/MinimalWordCountJava8.java b/examples/java8/src/main/java/org/apache/beam/examples/MinimalWordCountJava8.java index 24dd6f9..738b64d 100644 --- a/examples/java8/src/main/java/org/apache/beam/examples/MinimalWordCountJava8.java +++ b/examples/java8/src/main/java/org/apache/beam/examples/MinimalWordCountJava8.java @@ -67,6 +67,6 @@ public class MinimalWordCountJava8 { // CHANGE 3/3: The Google Cloud Storage path is required for outputting the results to. .apply(TextIO.Write.to("gs://YOUR_OUTPUT_BUCKET/AND_OUTPUT_PREFIX")); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/main/java/org/apache/beam/examples/complete/game/HourlyTeamScore.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/main/java/org/apache/beam/examples/complete/game/HourlyTeamScore.java b/examples/java8/src/main/java/org/apache/beam/examples/complete/game/HourlyTeamScore.java index aefa3bc..3a8d2ad 100644 --- a/examples/java8/src/main/java/org/apache/beam/examples/complete/game/HourlyTeamScore.java +++ b/examples/java8/src/main/java/org/apache/beam/examples/complete/game/HourlyTeamScore.java @@ -191,7 +191,7 @@ public class HourlyTeamScore extends UserScore { configureWindowedTableWrite())); - pipeline.run(); + pipeline.run().waitUntilFinish(); } // [END DocInclude_HTSMain] http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/main/java/org/apache/beam/examples/complete/game/UserScore.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/main/java/org/apache/beam/examples/complete/game/UserScore.java b/examples/java8/src/main/java/org/apache/beam/examples/complete/game/UserScore.java index f70b79c..32c939f 100644 --- a/examples/java8/src/main/java/org/apache/beam/examples/complete/game/UserScore.java +++ b/examples/java8/src/main/java/org/apache/beam/examples/complete/game/UserScore.java @@ -236,7 +236,7 @@ public class UserScore { configureBigQueryWrite())); // Run the batch pipeline. - pipeline.run(); + pipeline.run().waitUntilFinish(); } // [END DocInclude_USMain] http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/test/java/org/apache/beam/examples/complete/game/GameStatsTest.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/GameStatsTest.java b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/GameStatsTest.java index 7cd03f3..51ca719 100644 --- a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/GameStatsTest.java +++ b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/GameStatsTest.java @@ -69,7 +69,7 @@ public class GameStatsTest implements Serializable { // Check the set of spammers. PAssert.that(output).containsInAnyOrder(SPAMMERS); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/test/java/org/apache/beam/examples/complete/game/HourlyTeamScoreTest.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/HourlyTeamScoreTest.java b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/HourlyTeamScoreTest.java index f9fefb6..645f123 100644 --- a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/HourlyTeamScoreTest.java +++ b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/HourlyTeamScoreTest.java @@ -105,7 +105,7 @@ public class HourlyTeamScoreTest implements Serializable { PAssert.that(output).containsInAnyOrder(FILTERED_EVENTS); - p.run(); + p.run().waitUntilFinish(); } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/test/java/org/apache/beam/examples/complete/game/LeaderBoardTest.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/LeaderBoardTest.java b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/LeaderBoardTest.java index 9cba704..676dedb 100644 --- a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/LeaderBoardTest.java +++ b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/LeaderBoardTest.java @@ -110,7 +110,7 @@ public class LeaderBoardTest implements Serializable { .inOnTimePane(new IntervalWindow(baseTime, TEAM_WINDOW_DURATION)) .containsInAnyOrder(KV.of(blueTeam, 12), KV.of(redTeam, 4)); - p.run(); + p.run().waitUntilFinish(); } /** @@ -160,7 +160,7 @@ public class LeaderBoardTest implements Serializable { .inOnTimePane(window) .containsInAnyOrder(KV.of(blueTeam, 10), KV.of(redTeam, 9)); - p.run(); + p.run().waitUntilFinish(); } /** @@ -197,7 +197,7 @@ public class LeaderBoardTest implements Serializable { .inOnTimePane(window) .containsInAnyOrder(KV.of(redTeam, 14), KV.of(blueTeam, 13)); - p.run(); + p.run().waitUntilFinish(); } /** @@ -258,7 +258,7 @@ public class LeaderBoardTest implements Serializable { // account in earlier panes PAssert.that(teamScores).inFinalPane(window).containsInAnyOrder(KV.of(redTeam, 27)); - p.run(); + p.run().waitUntilFinish(); } /** @@ -346,7 +346,7 @@ public class LeaderBoardTest implements Serializable { KV.of(TestUser.BLUE_TWO.getUser(), 3), KV.of(TestUser.BLUE_TWO.getUser(), 8)); - p.run(); + p.run().waitUntilFinish(); } private TimestampedValue<GameActionInfo> event( http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e5afbe56/examples/java8/src/test/java/org/apache/beam/examples/complete/game/UserScoreTest.java ---------------------------------------------------------------------- diff --git a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/UserScoreTest.java b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/UserScoreTest.java index 7c86adf..39de333 100644 --- a/examples/java8/src/test/java/org/apache/beam/examples/complete/game/UserScoreTest.java +++ b/examples/java8/src/test/java/org/apache/beam/examples/complete/game/UserScoreTest.java @@ -110,7 +110,7 @@ public class UserScoreTest implements Serializable { // Check the user score sums. PAssert.that(output).containsInAnyOrder(USER_SUMS); - p.run(); + p.run().waitUntilFinish(); } /** Tests ExtractAndSumScore("team"). */ @@ -129,7 +129,7 @@ public class UserScoreTest implements Serializable { // Check the team score sums. PAssert.that(output).containsInAnyOrder(TEAM_SUMS); - p.run(); + p.run().waitUntilFinish(); } /** Test that bad input data is dropped appropriately. */ @@ -149,6 +149,6 @@ public class UserScoreTest implements Serializable { PAssert.that(extract).empty(); - p.run(); + p.run().waitUntilFinish(); } }