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)