MartijnVisser commented on code in PR #29399:
URL: https://github.com/apache/flink/pull/29399#discussion_r4196056692


##########
flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobVertexBackPressureHandlerTest.java:
##########
@@ -427,4 +432,86 @@ public long getLastUpdateTime() {
                                 .collect(Collectors.toList()))
                 .containsExactly(0, 1);
     }
+
+    /**
+     * The handler must not throw when a subtask's attempts are pruned 
concurrently (by {@code
+     * SubtaskMetricStore.retainAttempts()}) between its {@code size()}/{@code 
containsKey()} checks
+     * and the {@code keySet().iterator().next()} read. See FLINK-36146.
+     */
+    @Test
+    void testGetBackPressureWhenAttemptsRemovedConcurrently() throws Exception 
{
+        MetricStore multipleAttemptsMetricStore = new MetricStore();
+        for (MetricDump metricDump : getMultipleAttemptsMetricDumps()) {
+            multipleAttemptsMetricStore.add(metricDump);
+        }
+
+        // Empty subtask 0's attempts in the window between the handler's 
check and its iteration.
+        SubtaskMetricStore subtask =
+                multipleAttemptsMetricStore
+                        .getTaskMetricStore(
+                                
TEST_JOB_ID_BACK_PRESSURE_STATS_AVAILABLE.toString(),
+                                TEST_JOB_VERTEX_ID.toString())
+                        .getAllSubtaskMetricStores()
+                        .get(0);
+        Field attemptsField = 
SubtaskMetricStore.class.getDeclaredField("attempts");
+        attemptsField.setAccessible(true);
+        @SuppressWarnings("unchecked")
+        Map<Integer, ComponentMetricStore> originalAttempts =
+                (Map<Integer, ComponentMetricStore>) 
attemptsField.get(subtask);
+        attemptsField.set(subtask, new SelfEmptyingOnKeySet(originalAttempts));
+
+        JobVertexBackPressureHandler handler =
+                new JobVertexBackPressureHandler(
+                        () -> 
CompletableFuture.completedFuture(restfulGateway),
+                        Duration.ofSeconds(10),
+                        Collections.emptyMap(),
+                        JobVertexBackPressureHeaders.getInstance(),
+                        new MetricFetcher() {
+                            @Override
+                            public MetricStore getMetricStore() {
+                                return multipleAttemptsMetricStore;
+                            }
+
+                            @Override
+                            public void update() {}
+
+                            @Override
+                            public long getLastUpdateTime() {
+                                return 0;
+                            }
+                        });
+
+        Map<String, String> pathParameters = new HashMap<>();
+        pathParameters.put(
+                JobIDPathParameter.KEY, 
TEST_JOB_ID_BACK_PRESSURE_STATS_AVAILABLE.toString());
+        pathParameters.put(JobVertexIdPathParameter.KEY, 
TEST_JOB_VERTEX_ID.toString());
+        HandlerRequest<EmptyRequestBody> request =
+                HandlerRequest.resolveParametersAndCreate(
+                        EmptyRequestBody.getInstance(),
+                        new JobVertexMessageParameters(),
+                        pathParameters,
+                        Collections.emptyMap(),
+                        Collections.emptyList());
+
+        assertThatCode(() -> handler.handleRequest(request, restfulGateway))
+                .doesNotThrowAnyException();
+    }
+
+    private static final class SelfEmptyingOnKeySet
+            extends ConcurrentHashMap<Integer, ComponentMetricStore> {
+        private boolean emptied;
+
+        private SelfEmptyingOnKeySet(Map<Integer, ComponentMetricStore> 
initial) {
+            super(initial);
+        }
+
+        @Override
+        public KeySetView<Integer, ComponentMetricStore> keySet() {

Review Comment:
   With the copy, `keySet()` is never called on the live map, so this hook 
doesn't fire on this branch. Could it also clear on `entrySet()`? That still 
fails on master and passes here.



##########
flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobVertexBackPressureHandlerTest.java:
##########
@@ -427,4 +432,86 @@ public long getLastUpdateTime() {
                                 .collect(Collectors.toList()))
                 .containsExactly(0, 1);
     }
+
+    /**
+     * The handler must not throw when a subtask's attempts are pruned 
concurrently (by {@code
+     * SubtaskMetricStore.retainAttempts()}) between its {@code size()}/{@code 
containsKey()} checks
+     * and the {@code keySet().iterator().next()} read. See FLINK-36146.
+     */
+    @Test
+    void testGetBackPressureWhenAttemptsRemovedConcurrently() throws Exception 
{
+        MetricStore multipleAttemptsMetricStore = new MetricStore();
+        for (MetricDump metricDump : getMultipleAttemptsMetricDumps()) {
+            multipleAttemptsMetricStore.add(metricDump);
+        }
+
+        // Empty subtask 0's attempts in the window between the handler's 
check and its iteration.
+        SubtaskMetricStore subtask =
+                multipleAttemptsMetricStore
+                        .getTaskMetricStore(
+                                
TEST_JOB_ID_BACK_PRESSURE_STATS_AVAILABLE.toString(),
+                                TEST_JOB_VERTEX_ID.toString())
+                        .getAllSubtaskMetricStores()
+                        .get(0);
+        Field attemptsField = 
SubtaskMetricStore.class.getDeclaredField("attempts");
+        attemptsField.setAccessible(true);
+        @SuppressWarnings("unchecked")

Review Comment:
   Passing `subtask.getAllAttemptsMetricStores()` to the constructor avoids the 
cast and the `@SuppressWarnings`.



##########
flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobVertexBackPressureHandlerTest.java:
##########
@@ -427,4 +432,86 @@ public long getLastUpdateTime() {
                                 .collect(Collectors.toList()))
                 .containsExactly(0, 1);
     }
+
+    /**
+     * The handler must not throw when a subtask's attempts are pruned 
concurrently (by {@code
+     * SubtaskMetricStore.retainAttempts()}) between its {@code size()}/{@code 
containsKey()} checks
+     * and the {@code keySet().iterator().next()} read. See FLINK-36146.

Review Comment:
   I'd drop the ticket reference here, `git blame` leads to FLINK-40925 anyway.



-- 
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.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to