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();
   }

Reply via email to