Fix combine tests with Accumulation Mode These tests were not written in such a way as to succeed if the trigger fired multiple times.
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/e302cab0 Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/e302cab0 Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/e302cab0 Branch: refs/heads/master Commit: e302cab06640e016fcb24e03376ac84149f2ed2b Parents: f9fac64 Author: Thomas Groh <[email protected]> Authored: Wed Jul 27 09:04:06 2016 -0700 Committer: Kenneth Knowles <[email protected]> Committed: Wed Aug 24 12:46:24 2016 -0700 ---------------------------------------------------------------------- .../java/org/apache/beam/sdk/transforms/CombineTest.java | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/e302cab0/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/CombineTest.java ---------------------------------------------------------------------- diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/CombineTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/CombineTest.java index 6421b3b..897d17a 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/CombineTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/CombineTest.java @@ -21,7 +21,6 @@ import static org.apache.beam.sdk.TestUtils.checkCombineFn; import static org.apache.beam.sdk.transforms.display.DisplayDataMatchers.hasDisplayItem; import static org.apache.beam.sdk.transforms.display.DisplayDataMatchers.hasNamespace; import static org.apache.beam.sdk.transforms.display.DisplayDataMatchers.includesDisplayDataFrom; - import static com.google.common.base.Preconditions.checkArgument; import static com.google.common.base.Preconditions.checkNotNull; import static com.google.common.base.Preconditions.checkState; @@ -387,7 +386,7 @@ public class CombineTest implements Serializable { PCollection<String> output = input .apply(Window.<Integer>into(new GlobalWindows()) - .triggering(AfterPane.elementCountAtLeast(1)) + .triggering(Repeatedly.forever(AfterPane.elementCountAtLeast(1))) .accumulatingFiredPanes() .withAllowedLateness(new Duration(0))) .apply(Sum.integersGlobally()) @@ -583,7 +582,13 @@ public class CombineTest implements Serializable { .apply(Sum.integersGlobally().withoutDefaults().withFanout(2)) .apply(ParDo.of(new GetLast())); - PAssert.that(output).containsInAnyOrder(15); + PAssert.that(output).satisfies(new SerializableFunction<Iterable<Integer>, Void>() { + @Override + public Void apply(Iterable<Integer> input) { + assertThat(input, hasItem(15)); + return null; + } + }); pipeline.run(); }
