From cc7c24f198c82719a3356efebee6a50e0f0b5067 Mon Sep 17 00:00:00 2001 From: Yi Hu Date: Mon, 20 Jul 2026 13:23:19 -0400 Subject: [PATCH 1/2] Fix gRPC stream observer leak on ProcessBundleHandler shutdown * tears down active gRPC multiplexers on ProcessBundleHandler shutdown --- .../fn/harness/control/ProcessBundleHandler.java | 1 + .../beam/fn/harness/data/BeamFnDataClient.java | 9 ++++++++- .../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 449afd6a0243..b744da10cf6f 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 BeamFnApi.InstructionResponse.Builder trySplit(InstructionRequest request /** 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 1a50f5b448c5..458ed1209f43 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 @@ * 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 @@ void unregisterReceiver( /** Get the outbound observer for the specified apiServiceDescriptor and dataStreamId. */ StreamObserver 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 2f2a6b0fc660..8d3d7f76c18a 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 void poisonInstructionId(String instructionId) { } } + @Override + public void close() { + for (BeamFnDataGrpcMultiplexer client : multiplexerCache.values()) { + try { + client.close(); + } catch (Exception e) { + // proceed to close other clients + } + } + multiplexerCache.clear(); + } + @Override public StreamObserver getOutboundObserver( ApiServiceDescriptor apiServiceDescriptor, String dataStreamId) { From eacf9b91523c7e9041898d4a057e1ec05775bbf9 Mon Sep 17 00:00:00 2001 From: Yi Hu Date: Tue, 21 Jul 2026 10:02:05 -0400 Subject: [PATCH 2/2] 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> --- .../org/apache/beam/fn/harness/data/BeamFnDataGrpcClient.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 8d3d7f76c18a..79e34fe5765f 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 @@ -124,7 +124,7 @@ public void close() { try { client.close(); } catch (Exception e) { - // proceed to close other clients + LOG.warn("Failed to close multiplexer", e); } } multiplexerCache.clear();