[ 
https://issues.apache.org/jira/browse/FLINK-40644?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40644:
-----------------------------------
    Labels: pull-request-available  (was: )

> ForSt file cache: CachedDataInputStream.close() does not close the remote 
> (original) stream, leaking S3A connections until the pool is exhausted
> ------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40644
>                 URL: https://issues.apache.org/jira/browse/FLINK-40644
>             Project: Flink
>          Issue Type: Bug
>    Affects Versions: 2.3.0, 2.2.1, 2.1.3
>            Reporter: seungjoo choi
>            Priority: Critical
>              Labels: pull-request-available
>
> h3. Symptom
> A job using the ForSt state backend with the file cache and an S3A primary 
> directory runs normally for hours, then throughput drops to zero. Thread 
> dumps show the ForSt read threads blocked in 
> {{org.apache.http.pool.AbstractConnPool.getPoolEntryBlocking}}; after 
> {{s3.connection.timeout}} (~200 s) the job fails with 
> {{org.apache.hadoop.fs.s3a.AWSClientIOException: ... Timeout waiting for 
> connection from pool}}, restarts (which resets the pool), and the cycle 
> repeats. Raising {{s3.connection.maximum}} only delays the point of failure.
> h3. Root cause
> {{FileCacheEntry.open()}} opens the remote stream ({{FileSystem.open()}} on 
> the primary path), wraps it in a {{CachedDataInputStream}} and keeps no other 
> reference to it. {{CachedDataInputStream.close()}} only calls 
> {{closeCachedStream()}} (the local cache copy) and never closes 
> {{originalStream}}. Every stream that has read from S3 at least once holds a 
> leased HTTP connection in the S3A connection pool; the pool is not GC-aware, 
> so the lease is leaked for the lifetime of the TaskManager.
> The leak occurs whenever a stream is closed while its file is still cached: 
> table-cache eviction of an SST ({{state.backend.forst.files.open}} smaller 
> than the number of live files), deletion of compaction input files, and 
> overflow of the per-file stream pool in 
> {{ByteBufferReadableFSDataInputStream}}. The smaller {{files.open}}, the 
> faster the pool is exhausted (~58 reopens/s per TaskManager observed with 
> {{files.open=64}}).
> A related path: {{ByteBufferReadableFSDataInputStream.readFully(long, 
> ByteBuffer)}} neither closes nor returns the pooled stream when the read 
> throws.
> h3. Reproduction
> Flink 2.1.3, ForSt with {{state.backend.forst.cache.size-based-limit}} 
> smaller than the state (cache misses), S3A primary directory, keyed job with 
> a 1-day TTL state (~40 GB). With {{s3.connection.maximum=96}} the job hung 
> after ~6 hours, with 1024 after ~2 days. On the TaskManager, 
> {{/proc/net/tcp}} shows the number of ESTABLISHED connections to S3 growing 
> monotonically with the number of closed streams. With the fix it stays in the 
> single digits across hundreds of thousands of cache misses.
> h3. Fix
> In {{CachedDataInputStream.close()}} also close {{originalStream}} (after the 
> cached stream; propagate the first exception and attach the second as 
> suppressed), and close the pooled stream on the exception path of 
> {{ByteBufferReadableFSDataInputStream.readFully}}. No state-format or 
> configuration change.
> We have been running this fix in production as a class override on 2.1.3. I'd 
> like to work on this — could someone assign the ticket to me?



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

Reply via email to