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

Reply via email to