From b26675e6a51e3309cf54687b5ca3f5ff7cf4801f Mon Sep 17 00:00:00 2001 From: Rui Fan <1996fanrui@gmail.com> Date: Tue, 15 Sep 2026 12:43:35 +0200 Subject: [PATCH] [FLINK-40667][network] Fix buffer ref-count leak in recovery channel-state filtering --- .../channel/ChannelStateFilteringHandler.java | 2 ++ .../channel/RecoveredChannelStateHandler.java | 6 ++++- ...InputChannelRecoveredStateHandlerTest.java | 26 ++++++++++++++++--- 3 files changed, 30 insertions(+), 4 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/ChannelStateFilteringHandler.java b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/ChannelStateFilteringHandler.java index be9c91e495514..2151a3ea7f557 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/ChannelStateFilteringHandler.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/ChannelStateFilteringHandler.java @@ -118,6 +118,7 @@ public void filterAndRewrite( throws IOException { if (gateIndex < 0 || gateIndex >= gateHandlers.length) { + sourceBuffer.recycleBuffer(); throw new IllegalStateException( "Invalid gateIndex: " + gateIndex @@ -127,6 +128,7 @@ public void filterAndRewrite( GateFilterHandler gateHandler = gateHandlers[gateIndex]; if (gateHandler == null) { + sourceBuffer.recycleBuffer(); throw new IllegalStateException( "No handler for gateIndex " + gateIndex diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/RecoveredChannelStateHandler.java b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/RecoveredChannelStateHandler.java index 427f1d5cc9e3a..76e2677a46a74 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/RecoveredChannelStateHandler.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/channel/RecoveredChannelStateHandler.java @@ -588,12 +588,16 @@ public void recover( Buffer buffer = bufferWithContext.context; try { if (buffer.readableBytes() > 0) { + // Resolve the target serializer before retaining, so a failure here cannot leak + // the retained buffer reference. + DataOutputSerializer serializer = + segmentSerializerFor(getMappedChannels(channelInfo).getChannelInfo()); filteringHandler.filterAndRewrite( channelInfo.getGateIdx(), oldSubtaskIndex, channelInfo.getInputChannelIdx(), buffer.retainBuffer(), - segmentSerializerFor(getMappedChannels(channelInfo).getChannelInfo())); + serializer); } } finally { buffer.recycleBuffer(); diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/channel/InputChannelRecoveredStateHandlerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/channel/InputChannelRecoveredStateHandlerTest.java index 44050135904c8..a809ece62557f 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/channel/InputChannelRecoveredStateHandlerTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/channel/InputChannelRecoveredStateHandlerTest.java @@ -144,11 +144,11 @@ private AbstractInputChannelRecoveredStateHandler buildMultiChannelHandler() { /** Builds a handler in filtering mode (non-null filtering handler, no-op stub). */ private SpillingWithFilteringHandler buildFilteringInputChannelStateHandler() { - // Empty GateFilterHandler array: filtering is "enabled" structurally, but no gate-level - // filter logic runs. Suitable for exercising getBuffer() routing only. + // A single null gate handler: getBuffer routing is unaffected, and recover() hits the + // dispatcher's reachable null-handler throw path. ChannelStateFilteringHandler stubFilteringHandler = new ChannelStateFilteringHandler( - new ChannelStateFilteringHandler.GateFilterHandler[0]); + new ChannelStateFilteringHandler.GateFilterHandler[] {null}); return (SpillingWithFilteringHandler) AbstractInputChannelRecoveredStateHandler.create( new InputGate[] {inputGate}, @@ -351,6 +351,26 @@ void testPreFilterSegmentFreedOnClose() throws Exception { assertThat(filteringHandler.getPreFilterSegmentForTesting()).isNull(); } + @Test + void testPreFilterBufferRecycledWhenFilterAndRewriteThrows() throws Exception { + // On the dispatcher's null-handler throw, the pre-filter buffer must still be recycled, or + // close() would free() a segment still wrapped by a live NetworkBuffer. + try (SpillingWithFilteringHandler filteringHandler = + buildFilteringInputChannelStateHandler()) { + RecoveredChannelStateHandler.BufferWithContext bwc = + filteringHandler.getBuffer(channelInfo); + // Non-empty buffer so recover() reaches filterAndRewrite. + bwc.context.setSize(Long.BYTES); + assertThat(filteringHandler.isPreFilterBufferInUse()).isTrue(); + + assertThatThrownBy(() -> filteringHandler.recover(channelInfo, 0, bwc)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("No handler for gateIndex"); + + assertThat(filteringHandler.isPreFilterBufferInUse()).isFalse(); + } + } + @Test void testSpillingHandlerRequiresSpillDirectories() { assertThatThrownBy(() -> buildSpillingNoFilteringHandler(null))