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..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 @@ -118,6 +118,18 @@ public void poisonInstructionId(String instructionId) { } } + @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 getOutboundObserver( ApiServiceDescriptor apiServiceDescriptor, String dataStreamId) {