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