zhijiangW commented on a change in pull request #9170: [FLINK-13325][test] Add 
test to guard the fix for FLINK-13249
URL: https://github.com/apache/flink/pull/9170#discussion_r306182783
 
 

 ##########
 File path: 
flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/consumer/RemoteInputChannelTest.java
 ##########
 @@ -1081,6 +1094,79 @@ private RemoteInputChannel createRemoteInputChannel(
                        .buildRemoteAndSetToGate(inputGate);
        }
 
+       /**
+        * Test to guard against FLINK-13249.
+        */
+       @Test
+       public void testOnFailedPartitionRequestDoesNotBlockNetworkThreads() 
throws Exception {
+
+               final long testBlockedWaitTimeoutMillis = 30_000L;
+
+               final PartitionProducerStateChecker 
partitionProducerStateChecker =
+                       (jobId, intermediateDataSetId, resultPartitionId) -> 
CompletableFuture.completedFuture(ExecutionState.RUNNING);
+               final NettyShuffleEnvironment shuffleEnvironment = new 
NettyShuffleEnvironmentBuilder().build();
+               final Task task = new TestTaskBuilder(shuffleEnvironment)
+                       
.setPartitionProducerStateChecker(partitionProducerStateChecker)
+                       .build();
+               final SingleInputGate inputGate = new SingleInputGateBuilder()
+                       .setPartitionProducerStateProvider(task)
+                       .build();
+
+               TestTaskBuilder.setTaskState(task, ExecutionState.RUNNING);
+
+               final OneShotLatch ready = new OneShotLatch();
+               final OneShotLatch blocker = new OneShotLatch();
+               final AtomicBoolean timedOutOrInterrupted = new 
AtomicBoolean(false);
+
+               final ConnectionManager blockingConnectionManager = new 
TestingConnectionManager() {
+
+                       @Override
+                       public PartitionRequestClient 
createPartitionRequestClient(
+                               ConnectionID connectionId) {
+                               ready.trigger();
+                               try {
+                                       // We block here, in a section that 
holds the SingleInputGate#requestLock
+                                       
blocker.await(testBlockedWaitTimeoutMillis, TimeUnit.MILLISECONDS);
+                               } catch (InterruptedException | 
TimeoutException e) {
+                                       timedOutOrInterrupted.set(true);
+                               }
+
+                               return new TestingPartitionRequestClient();
+                       }
+               };
+
+               final RemoteInputChannel remoteInputChannel =
+                       InputChannelBuilder.newBuilder()
+                               .setConnectionManager(blockingConnectionManager)
+                               .buildRemoteAndSetToGate(inputGate);
+
+               final Thread simulatedNetworkThread = new Thread(
+                       () -> {
+                               try {
+                                       ready.await();
+                                       // We want to make sure that our 
simulated network thread does not block on
+                                       // SingleInputGate#requestLock as well 
through this call.
+                                       
remoteInputChannel.onFailedPartitionRequest();
+
+                                       // Will only give free the blocker if 
we did not block ourselves.
+                                       blocker.trigger();
+                               } catch (InterruptedException e) {
+                                       Thread.currentThread().interrupt();
+                               }
+                       });
+
+               simulatedNetworkThread.start();
+
+               // The entry point to that will lead us into 
blockingConnectionManager#createPartitionRequestClient(...).
+               inputGate.requestPartitions();
+
+               simulatedNetworkThread.join();
+
+               Assert.assertFalse(
+                       "Test ended by timeout or interruption - this indicates 
that the network thread was blocked.",
+                       timedOutOrInterrupted.get());
+       }
 
 Review comment:
   Might be better to close `NettyShuffleEnvironment` and `SingleInputGate` 
finally.

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


With regards,
Apache Git Services

Reply via email to