This is an automated email from the ASF dual-hosted git repository.
stankiewicz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 9c82e053f02 Avoid output inside try-catch in Java IO (#39124)
9c82e053f02 is described below
commit 9c82e053f02ca154cb9af91bf71b90d93a02ee45
Author: Aryankn29 <[email protected]>
AuthorDate: Tue Jul 28 03:53:08 2026 -0700
Avoid output inside try-catch in Java IO (#39124)
* Avoid output inside try-catch in Java IO
* Preserve null handling when moving output
* Apply suggestion from @gemini-code-assist[bot]
Co-authored-by: gemini-code-assist[bot]
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
* Apply Spotless formatting
---------
Co-authored-by: gemini-code-assist[bot]
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
---
.../org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java | 17 +++++++++++++++--
.../org/apache/beam/sdk/io/gcp/healthcare/FhirIO.java | 6 +++++-
.../org/apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java | 17 +++++++++++++----
3 files changed, 33 insertions(+), 7 deletions(-)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java
index b222b358f54..c2ac4efb5b5 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.java
@@ -2105,9 +2105,13 @@ public class BigQueryIO {
// the same order.
BoundedSource.BoundedReader<T> reader =
streamSource.createReader(options);
+ T current = null;
+ boolean hasCurrent = false;
try {
if (reader.start()) {
- outputReceiver.get(rowTag).output(reader.getCurrent());
+ current =
+ java.util.Objects.requireNonNull(reader.getCurrent(), "Reader
returned null element");
+ hasCurrent = true;
} else {
return;
}
@@ -2120,11 +2124,17 @@ public class BigQueryIO {
(Exception) e.getCause(),
"Unable to parse record reading from BigQuery");
}
+ if (hasCurrent) {
+ outputReceiver.get(rowTag).output(current);
+ }
while (true) {
+ current = null;
+ hasCurrent = false;
try {
if (reader.advance()) {
- outputReceiver.get(rowTag).output(reader.getCurrent());
+ current = reader.getCurrent();
+ hasCurrent = true;
} else {
return;
}
@@ -2137,6 +2147,9 @@ public class BigQueryIO {
(Exception) e.getCause(),
"Unable to parse record reading from BigQuery");
}
+ if (hasCurrent) {
+ outputReceiver.get(rowTag).output(current);
+ }
}
}
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/FhirIO.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/FhirIO.java
index 70be676c75f..1162db9e1a2 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/FhirIO.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/FhirIO.java
@@ -602,8 +602,9 @@ public class FhirIO {
@ProcessElement
public void processElement(ProcessContext context) {
String resourceId = context.element();
+ String resource = null;
try {
- context.output(fetchResource(this.client, resourceId));
+ resource =
java.util.Objects.requireNonNull(fetchResource(this.client, resourceId));
} catch (Exception e) {
READ_RESOURCE_ERRORS.inc();
LOG.warn(
@@ -612,6 +613,9 @@ public class FhirIO {
e);
context.output(FhirIO.Read.DEAD_LETTER,
HealthcareIOError.of(resourceId, e));
}
+ if (resource != null) {
+ context.output(resource);
+ }
}
private String fetchResource(HealthcareApiClient client, String
resourceName)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java
index 3647ef7671e..496b39a923f 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/healthcare/HL7v2IO.java
@@ -365,11 +365,15 @@ public class HL7v2IO {
@ProcessElement
public void processElement(ProcessContext context) {
String msgId = context.element();
+ HL7v2Message message = null;
try {
- context.output(client.fetchMessage(msgId));
+ message =
java.util.Objects.requireNonNull(client.fetchMessage(msgId));
} catch (Exception e) {
context.output(HL7v2IO.Read.DEAD_LETTER,
HealthcareIOError.of(msgId, e));
}
+ if (message != null) {
+ context.output(message);
+ }
}
}
}
@@ -487,15 +491,20 @@ public class HL7v2IO {
@ProcessElement
public void processElement(ProcessContext context) {
String msgId = context.element().getHl7v2MessageId();
+ HL7v2ReadResponse response = null;
try {
- HL7v2ReadResponse response =
- HL7v2ReadResponse.of(context.element().getMetadata(),
client.fetchMessage(msgId));
- context.output(response);
+ response =
+ java.util.Objects.requireNonNull(
+ HL7v2ReadResponse.of(
+ context.element().getMetadata(),
client.fetchMessage(msgId)));
} catch (Exception e) {
HealthcareIOError<HL7v2ReadParameter> error =
HealthcareIOError.of(context.element(), e);
context.output(HL7v2IO.HL7v2Read.DEAD_LETTER, error);
}
+ if (response != null) {
+ context.output(response);
+ }
}
}
}