Martijn Visser created FLINK-40660:
--------------------------------------

             Summary: SplitFetcher hangs on shutdown when its element queue 
wakeUp flag is still set
                 Key: FLINK-40660
                 URL: https://issues.apache.org/jira/browse/FLINK-40660
             Project: Flink
          Issue Type: Bug
          Components: Connectors / Common
            Reporter: Martijn Visser


{{SplitFetcher#run}} puts an empty synchronisation batch into the element queue 
after its loop ends
and discards the return value:

{code:java}
elementsQueue.put(
        fetcherId(),
        new RecordsBySplits<E>(Collections.emptyMap(), Collections.emptySet()) {
{code}

{{put}} returns false without enqueueing when the queue is full and that 
fetcher's wakeUp flag is
set. The batch is then dropped, {{recycle()}} never runs, 
{{recordsProcessedLatch.await()}} never
returns, the {{SplitReader}} is never closed and the shutdown hook never runs, 
so
{{fetchersToShutDown}} is never decremented and {{SplitFetcherManager#close}} 
blocks until its
timeout. {{maybeShutdownFinishedFetchers}} has already removed the fetcher from 
the map, so nothing
else can release it.

Reproduced on master by filling a capacity-1 queue, calling 
{{wakeUpPuttingThread}} for the fetcher's
own index and then {{shutdown(true)}}: the shutdown hook does not run, the 
reader is not closed and
the fetcher thread sits in WAITING. The flag is sticky, and the {{finally}} in 
{{FetchTask#run}}
clears only its own field, so a {{wakeUp}} that arrives just after the records 
were enqueued leaves
it set with nothing to consume it.

The put should be retried rather than dropped.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to