From 87dc61f927c0dcd7815f080aa16b06c1403c3e78 Mon Sep 17 00:00:00 2001 From: Utkarsh Date: Tue, 23 Jun 2026 21:02:31 -0700 Subject: [PATCH 1/4] Fix JsonToRow swallowing downstream errors when runners fuse transforms. Separate JSON parsing from MultiOutputReceiver output in ParseWithError so exceptions from fused downstream consumers are not misreported as parse failures. Fixes #20935. --- .../apache/beam/sdk/transforms/JsonToRow.java | 8 +++--- .../beam/sdk/transforms/JsonToRowTest.java | 25 +++++++++++++++++++ 2 files changed, 29 insertions(+), 4 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/JsonToRow.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/JsonToRow.java index d812d299a836..160dfd3b43bf 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/JsonToRow.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/JsonToRow.java @@ -309,12 +309,10 @@ public static ParseWithError create(JsonToRowWithErrFn jsonToRowWithErrFn) { @ProcessElement public void processElement(@Element String element, MultiOutputReceiver output) { + final Row parsedRow; try { - - output.get(PARSED_LINE).output(jsonToRow(objectMapper(), element)); - + parsedRow = jsonToRow(objectMapper(), element); } catch (Exception ex) { - if (getJsonToRowWithErrFn().getExtendedErrorInfo()) { output .get(PARSE_ERROR) @@ -328,7 +326,9 @@ public void processElement(@Element String element, MultiOutputReceiver output) .get(PARSE_ERROR) .output(Row.withSchema(ERROR_ROW_SCHEMA).addValue(element).build()); } + return; } + output.get(PARSED_LINE).output(parsedRow); } private ObjectMapper objectMapper() { diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java index 490cb68ab9e4..0e6408e32269 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java @@ -274,6 +274,31 @@ public void testParsesErrorWithErrorMsgWithRequireNullDeadLetter() throws Except pipeline.run(); } + @Test + @Category(NeedsRunner.class) + public void testDownstreamExceptionIsNotReportedAsParseError() { + PCollection jsonPersons = pipeline.apply("jsonPersons", Create.of(JSON_PERSON.get(0))); + + ParseResult results = jsonPersons.apply(JsonToRow.withExceptionReporting(PERSON_SCHEMA)); + + results + .getResults() + .apply( + "throwingDownstream", + ParDo.of( + new DoFn() { + @ProcessElement + public void processElement(ProcessContext context) { + throw new RuntimeException("downstream failure"); + } + })); + + thrown.expect(RuntimeException.class); + thrown.expectMessage("downstream failure"); + + pipeline.run(); + } + @Test @Category(NeedsRunner.class) public void testParsesErrorWithErrorMsgRowsDeadLetterWithCustomFieldNames() throws Exception { From 0bad579499bf684e4c89262e4de2de52b3d3e4d6 Mon Sep 17 00:00:00 2001 From: Utkarsh Date: Wed, 24 Jun 2026 19:54:07 -0700 Subject: [PATCH 2/4] Address review: use MapElements in JsonToRow regression test. Avoid anonymous DoFn capturing the non-serializable test instance so the test is safe on runners that enforce DoFn serialization. --- .../apache/beam/sdk/transforms/JsonToRowTest.java | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java index 0e6408e32269..80bb8f11e5d9 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java @@ -33,6 +33,7 @@ import org.apache.beam.sdk.util.RowJson.RowJsonDeserializer.NullBehavior; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.Row; +import org.apache.beam.sdk.values.TypeDescriptors; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.junit.Rule; import org.junit.Test; @@ -285,13 +286,11 @@ public void testDownstreamExceptionIsNotReportedAsParseError() { .getResults() .apply( "throwingDownstream", - ParDo.of( - new DoFn() { - @ProcessElement - public void processElement(ProcessContext context) { - throw new RuntimeException("downstream failure"); - } - })); + MapElements.into(TypeDescriptors.rows()) + .via( + (Row row) -> { + throw new RuntimeException("downstream failure"); + })); thrown.expect(RuntimeException.class); thrown.expectMessage("downstream failure"); From 284bec9793b53175d20f09a3ba0936740b36cc9a Mon Sep 17 00:00:00 2001 From: Utkarsh Date: Wed, 24 Jun 2026 20:28:49 -0700 Subject: [PATCH 3/4] Fix JsonToRow regression test for runner integration suites. Use a static DoFn to avoid serialization issues, set row schema on the downstream transform, and tag the test with ValidatesRunner so it can run via Dataflow validatesRunner tasks. Co-authored-by: Cursor --- .../beam/sdk/transforms/JsonToRowTest.java | 20 ++++++++++--------- 1 file changed, 11 insertions(+), 9 deletions(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java index 80bb8f11e5d9..590bba01be0e 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java @@ -27,13 +27,13 @@ import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.testing.UsesSchema; +import org.apache.beam.sdk.testing.ValidatesRunner; import org.apache.beam.sdk.transforms.JsonToRow.JsonToRowFn; import org.apache.beam.sdk.transforms.JsonToRow.JsonToRowWithErrFn; import org.apache.beam.sdk.transforms.JsonToRow.ParseResult; import org.apache.beam.sdk.util.RowJson.RowJsonDeserializer.NullBehavior; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.Row; -import org.apache.beam.sdk.values.TypeDescriptors; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.junit.Rule; import org.junit.Test; @@ -276,7 +276,7 @@ public void testParsesErrorWithErrorMsgWithRequireNullDeadLetter() throws Except } @Test - @Category(NeedsRunner.class) + @Category({NeedsRunner.class, ValidatesRunner.class}) public void testDownstreamExceptionIsNotReportedAsParseError() { PCollection jsonPersons = pipeline.apply("jsonPersons", Create.of(JSON_PERSON.get(0))); @@ -284,13 +284,8 @@ public void testDownstreamExceptionIsNotReportedAsParseError() { results .getResults() - .apply( - "throwingDownstream", - MapElements.into(TypeDescriptors.rows()) - .via( - (Row row) -> { - throw new RuntimeException("downstream failure"); - })); + .apply("throwingDownstream", ParDo.of(new ThrowingDownstreamDoFn())) + .setRowSchema(PERSON_SCHEMA); thrown.expect(RuntimeException.class); thrown.expectMessage("downstream failure"); @@ -354,4 +349,11 @@ private static String jsonPerson(String name, String height) { private static Row row(Schema schema, Object... values) { return Row.withSchema(schema).addValues(values).build(); } + + private static class ThrowingDownstreamDoFn extends DoFn { + @ProcessElement + public void processElement(ProcessContext context) { + throw new RuntimeException("downstream failure"); + } + } } From 94d90e27cffb48c20a55eefc41b997ad98827e79 Mon Sep 17 00:00:00 2001 From: Utkarsh Date: Thu, 25 Jun 2026 02:22:12 -0700 Subject: [PATCH 4/4] Merge apache/beam master and fix JsonToRowTest for CI. Rebase onto current master (1141 commits behind) and remove ValidatesRunner category so the regression test only runs on DirectRunner needsRunnerTests. Suppress UnusedVariable on ThrowingDownstreamDoFn for Error Prone. Co-authored-by: Cursor --- .../java/org/apache/beam/sdk/transforms/JsonToRowTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java index 590bba01be0e..79918930e098 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/JsonToRowTest.java @@ -27,7 +27,6 @@ import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.testing.UsesSchema; -import org.apache.beam.sdk.testing.ValidatesRunner; import org.apache.beam.sdk.transforms.JsonToRow.JsonToRowFn; import org.apache.beam.sdk.transforms.JsonToRow.JsonToRowWithErrFn; import org.apache.beam.sdk.transforms.JsonToRow.ParseResult; @@ -276,7 +275,7 @@ public void testParsesErrorWithErrorMsgWithRequireNullDeadLetter() throws Except } @Test - @Category({NeedsRunner.class, ValidatesRunner.class}) + @Category(NeedsRunner.class) public void testDownstreamExceptionIsNotReportedAsParseError() { PCollection jsonPersons = pipeline.apply("jsonPersons", Create.of(JSON_PERSON.get(0))); @@ -352,6 +351,7 @@ private static Row row(Schema schema, Object... values) { private static class ThrowingDownstreamDoFn extends DoFn { @ProcessElement + @SuppressWarnings("UnusedVariable") public void processElement(ProcessContext context) { throw new RuntimeException("downstream failure"); }