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