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) {

Reply via email to