This is an automated email from the ASF dual-hosted git repository.
Abacn 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 0c151d8ddb2 Fix gRPC stream observer leak on ProcessBundleHandler
shutdown (#39390)
0c151d8ddb2 is described below
commit 0c151d8ddb23df45b96d660bb8b7e644bece7068
Author: Yi Hu <[email protected]>
AuthorDate: Tue Jul 21 12:41:46 2026 -0400
Fix gRPC stream observer leak on ProcessBundleHandler shutdown (#39390)
* Fix gRPC stream observer leak on ProcessBundleHandler shutdown
* tears down active gRPC multiplexers on ProcessBundleHandler shutdown
* Update
sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java
Co-authored-by: gemini-code-assist[bot]
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
---------
Co-authored-by: gemini-code-assist[bot]
<176961590+gemini-code-assist[bot]@users.noreply.github.com>
---
.../apache/beam/fn/harness/control/ProcessBundleHandler.java | 1 +
.../org/apache/beam/fn/harness/data/BeamFnDataClient.java | 9 ++++++++-
.../apache/beam/fn/harness/data/BeamFnDataGrpcClient.java | 12 ++++++++++++
3 files changed, 21 insertions(+), 1 deletion(-)
diff --git
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/ProcessBundleHandler.java
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/ProcessBundleHandler.java
index 449afd6a024..b744da10cf6 100644
---
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/ProcessBundleHandler.java
+++
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/control/ProcessBundleHandler.java
@@ -767,6 +767,7 @@ public class ProcessBundleHandler {
/** Shutdown the bundles, running the tearDown() functions. */
public void shutdown() throws Exception {
bundleProcessorCache.shutdown();
+ beamFnDataClient.close();
}
@VisibleForTesting
diff --git
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataClient.java
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataClient.java
index 1a50f5b448c..458ed1209f4 100644
---
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataClient.java
+++
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataClient.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.fn.harness.data;
+import java.io.Closeable;
+import java.io.IOException;
import java.util.List;
import org.apache.beam.model.fnexecution.v1.BeamFnApi.Elements;
import org.apache.beam.model.pipeline.v1.Endpoints;
@@ -30,7 +32,7 @@ import
org.apache.beam.vendor.grpc.v1p69p0.io.grpc.stub.StreamObserver;
* provide a receiver of outbound elements. Callers can register themselves as
receivers for inbound
* elements or can get a handle for a receiver of outbound elements.
*/
-public interface BeamFnDataClient {
+public interface BeamFnDataClient extends Closeable {
/**
* Registers a receiver for the provided instruction id.
*
@@ -72,4 +74,9 @@ public interface BeamFnDataClient {
/** Get the outbound observer for the specified apiServiceDescriptor and
dataStreamId. */
StreamObserver<Elements> getOutboundObserver(
Endpoints.ApiServiceDescriptor apiServiceDescriptor, String
dataStreamId);
+
+ @Override
+ default void close() throws IOException {
+ // Default to no-op
+ }
}
diff --git
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java
index 2f2a6b0fc66..79e34fe5765 100644
---
a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java
+++
b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java
@@ -118,6 +118,18 @@ public class BeamFnDataGrpcClient implements
BeamFnDataClient {
}
}
+ @Override
+ public void close() {
+ for (BeamFnDataGrpcMultiplexer client : multiplexerCache.values()) {
+ try {
+ client.close();
+ } catch (Exception e) {
+ LOG.warn("Failed to close multiplexer", e);
+ }
+ }
+ multiplexerCache.clear();
+ }
+
@Override
public StreamObserver<Elements> getOutboundObserver(
ApiServiceDescriptor apiServiceDescriptor, String dataStreamId) {