snuyanzin commented on code in PR #29386:
URL: https://github.com/apache/flink/pull/29386#discussion_r4188325041


##########
flink-connectors/flink-connector-base/src/test/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManagerTest.java:
##########
@@ -309,6 +311,40 @@ void 
testIdleShutdownSplitFetcherWaitsUntilRecordProcessed() throws Exception {
         }
     }
 
+    @RepeatedTest(100)
+    void testAddSplitsWhileClosing() throws Exception {

Review Comment:
   ```java
   @Test
   void getRunningFetcherSurvivesFetcherRemovalInToctouWindow() throws 
Exception {
       final SingleThreadFetcherManager<Object, TestingSourceSplit> 
fetcherManager =
               new SingleThreadFetcherManager<>(TestingSplitReader::new, new 
Configuration());
   
       final SplitFetcher<Object, TestingSourceSplit> fetcher = 
fetcherManager.createSplitFetcher();
   
       final Map<Integer, SplitFetcher<Object, TestingSourceSplit>> 
racyFetchers =
               new ConcurrentHashMap<Integer, SplitFetcher<Object, 
TestingSourceSplit>>() {
                   private boolean removed;
   
                   @Override
                   public Collection<SplitFetcher<Object, TestingSourceSplit>> 
values() {
                       // The shutdown callback (fetchers.remove(...)) fires 
here: after
                       // getRunningFetcher()'s isEmpty() check saw the entry, 
before it is read.
                       if (!removed) {
                           removed = true;
                           clear();
                       }
                       return super.values();
                   }
               };
       racyFetchers.put(0, fetcher);
   
       final Field fetchersField = 
SplitFetcherManager.class.getDeclaredField("fetchers");
       fetchersField.setAccessible(true);
       fetchersField.set(fetcherManager, racyFetchers);
   
       assertThat(fetcherManager.getRunningFetcher()).isNull();
   }
   ```



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