diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/consumer/UnknownInputChannel.java b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/consumer/UnknownInputChannel.java index 9db9df9c613c14..c5a12bcb829bcf 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/consumer/UnknownInputChannel.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/consumer/UnknownInputChannel.java @@ -192,6 +192,8 @@ public RemoteInputChannel toRemoteInputChannel( metrics.getNumBytesInRemoteCounter(), metrics.getNumBuffersInRemoteCounter(), channelStateWriter == null ? ChannelStateWriter.NO_OP : channelStateWriter, + // Unknown channels exist only in BATCH jobs, which have no channel + // state, so this channel is never in recovery. false); return channel; } @@ -210,6 +212,8 @@ public LocalInputChannel toLocalInputChannel(ResultPartitionID resultPartitionID metrics.getNumBuffersInLocalCounter(), channelStateWriter == null ? ChannelStateWriter.NO_OP : channelStateWriter, networkBuffersPerChannel, + // Unknown channels exist only in BATCH jobs, which have no channel + // state, so this channel is never in recovery. false); }