[ 
https://issues.apache.org/jira/browse/FLINK-25672?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17941577#comment-17941577
 ] 

Nickel Fang commented on FLINK-25672:
-------------------------------------

We had the same issue (OOM in job manager) and have optimized it by TTL in our 
project. Can I request a PR?


h1. Solution

Requirements
 # Can delete the expired processed paths according to TTL policy to save the 
memory.
 # Be an enhancement solution. Users can easily choose to use the original 
solution or the enhancement solution
 # Can seamlessly upgrade to the enhancement solution with no data lost for the 
running streaming job.

Prerequisite

There is a TTL mechanism in the file source (e.g. S3 retention policy).


Details

Introduce a new variable LinkedHashSet<Pair<Path, Long>> 
alreadyProcessedPathAndTimestamps in 
PendingSplitsCheckpoint, and Duration retentionTime in 
ContinuousEnumerationSettings.
 

> FileSource enumerator remembers paths of all already processed files which 
> can result in large state
> ----------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-25672
>                 URL: https://issues.apache.org/jira/browse/FLINK-25672
>             Project: Flink
>          Issue Type: Improvement
>          Components: Connectors / FileSystem
>            Reporter: Martijn Visser
>            Priority: Major
>
> As mentioned in the Filesystem documentation, for Unbounded File Sources, the 
> {{FileEnumerator}} currently remembers paths of all already processed files, 
> which is a state that can in come cases grow rather large. 
> We should look into possibilities to reduce this. We could look into adding a 
> compressed form of tracking already processed files (for example by keeping 
> modification timestamps lower boundaries).
> When fixed, this should also be reflected in the documentation, as mentioned 
> in https://github.com/apache/flink/pull/18288#discussion_r785707311



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

Reply via email to