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