Don't Suppress Throwable in PAssert in Streaming Mode

Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/ea3a46c7
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/ea3a46c7
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/ea3a46c7

Branch: refs/heads/master
Commit: ea3a46c7896c954084fc581a06e1f3ef68a3f3d0
Parents: 4aba224
Author: Aljoscha Krettek <[email protected]>
Authored: Sun Jul 24 12:04:28 2016 +0200
Committer: Kenneth Knowles <[email protected]>
Committed: Wed Aug 24 12:46:24 2016 -0700

----------------------------------------------------------------------
 .../org/apache/beam/sdk/testing/PAssert.java    | 20 +++-----------------
 1 file changed, 3 insertions(+), 17 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/ea3a46c7/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
----------------------------------------------------------------------
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
index 3f1a741..943ed11 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/testing/PAssert.java
@@ -1066,14 +1066,7 @@ public class PAssert {
 
     @ProcessElement
     public void processElement(ProcessContext c) {
-      try {
-        doChecks(c.element(), checkerFn, success, failure);
-      } catch (Throwable t) {
-        // Suppress exception in streaming
-        if (!c.getPipelineOptions().as(StreamingOptions.class).isStreaming()) {
-          throw t;
-        }
-      }
+      doChecks(c.element(), checkerFn, success, failure);
     }
   }
 
@@ -1098,15 +1091,8 @@ public class PAssert {
 
     @ProcessElement
     public void processElement(ProcessContext c) {
-      try {
-        ActualT actualContents = Iterables.getOnlyElement(c.element());
-        doChecks(actualContents, checkerFn, success, failure);
-      } catch (Throwable t) {
-        // Suppress exception in streaming
-        if (!c.getPipelineOptions().as(StreamingOptions.class).isStreaming()) {
-          throw t;
-        }
-      }
+      ActualT actualContents = Iterables.getOnlyElement(c.element());
+      doChecks(actualContents, checkerFn, success, failure);
     }
   }
 

Reply via email to