[
https://issues.apache.org/jira/browse/FLINK-40659?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123454#comment-18123454
]
Martijn Visser commented on FLINK-40659:
----------------------------------------
With FLINK-40657 merged, a reaped fetcher no longer leaves state behind in the
element queue, so each cycle costs the same and nothing grows. What is left per
split is a {{SplitReader}} create and close, one empty batch through the queue
and five INFO lines; fetcher threads come from a cached pool and are reused.
Keeping an idle fetcher open needs a new option on {{SourceReaderOptions}} and
has to force the shutdown once {{noMoreSplitsAssignment}} is set, otherwise
bounded readers never reach END_OF_INPUT. Lowering this to Minor until there is
a measurement showing the remaining cost matters for a real job.
> SingleThreadFetcherManager reaps and recreates its fetcher on every idle gap
> ----------------------------------------------------------------------------
>
> Key: FLINK-40659
> URL: https://issues.apache.org/jira/browse/FLINK-40659
> Project: Flink
> Issue Type: Improvement
> Components: Connectors / Common
> Reporter: Martijn Visser
> Priority: Major
>
> {{SingleThreadFetcherManager#addSplits}} creates a new {{SplitFetcher}}
> whenever the fetcher map is
> empty, and {{SourceReaderBase#finishedOrAvailableLater}} calls
> {{maybeShutdownFinishedFetchers()}}
> every time the element queue drains, which reaps the fetcher as soon as it is
> idle. For a source that
> finishes a split and then requests the next one, such as the file source with
> continuous discovery,
> that is one fetcher per split.
> Each one costs a {{SplitReader}} construction and close, a pool thread, a
> {{FetchTask}}, a
> {{CountDownLatch}}, an extra empty batch through the element queue and five
> INFO log lines. The job
> reported in FLINK-40657 does this about 19 times a second per TaskManager for
> 17 hours, and there is
> no option to hold the fetcher open. It would be better to keep an idle
> fetcher for a short while than
> to reap it on every gap. FLINK-36146 reports a race in
> {{getRunningFetcher()}} caused by the same
> cycle.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)