[
https://issues.apache.org/jira/browse/FLINK-40657?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-40657:
-----------------------------------
Labels: pull-request-available (was: )
> FutureCompletingBlockingQueue retains a ConditionAndFlag for every
> SplitFetcher ever created
> ---------------------------------------------------------------------------------------------
>
> Key: FLINK-40657
> URL: https://issues.apache.org/jira/browse/FLINK-40657
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Common
> Affects Versions: 2.0.2, 2.3.0, 2.2.1, 2.1.3, 1.20.6
> Reporter: Alexandru
> Priority: Major
> Labels: pull-request-available
>
> {{FutureCompletingBlockingQueue}} keeps per-producer wakeup state in
> [an array indexed by a "thread
> index"|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java#L105-L107],
> [initialised to length
> 1|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java#L118]
> and grown on demand by
> [maybeCreateCondition|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java#L392-L401]:
> {code:java}
> // FutureCompletingBlockingQueue.java, L105-L107
> /** The per-thread conditions and wakeUp flags. */
> @GuardedBy("lock")
> private ConditionAndFlag[] putConditionAndFlags;
> // FutureCompletingBlockingQueue.java, L392-L401
> @GuardedBy("lock")
> private void maybeCreateCondition(int threadIndex) {
> if (putConditionAndFlags.length < threadIndex + 1) {
> putConditionAndFlags = Arrays.copyOf(putConditionAndFlags,
> threadIndex + 1);
> }
> if (putConditionAndFlags[threadIndex] == null) {
> putConditionAndFlags[threadIndex] = new
> ConditionAndFlag(lock.newCondition());
> }
> }
> {code}
> Nothing ever shrinks the array or nulls entries that are no longer needed,
> including on
> {{close()}}. Each populated slot pins a
> [ConditionAndFlag|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/synchronization/FutureCompletingBlockingQueue.java#L415]
> and the {{lock.newCondition()}} inside it.
> That would be fine if {{threadIndex}} were bounded by the number of I/O
> threads, as the field name
> and its javadoc imply. It is not. The index is the {{SplitFetcher}} id,
> passed in
> [when a fetcher hands a batch
> over|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcher.java#L137-L138],
> and that id comes from
> [a counter that is never
> reused|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java#L246-L253].
> The same index also reaches the queue via
> [FetchTask|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/FetchTask.java#L64]
> ({{elementsQueue.put(fetcherIndex, lastRecords)}}), where {{fetcherIndex}} is
> the same
> {{SplitFetcher}} id -- the task is constructed with {{new FetchTask<>(...,
> id)}}. On 1.x that is
> the *only* path, since {{SplitFetcher}} there has no direct {{put}} call; the
> line is byte-identical
> on 1.20 and master, so the leak is the same on both lines.
> Fetchers are short-lived.
> [SourceReaderBase#finishedOrAvailableLater|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/SourceReaderBase.java#L433-L434]
> calls
> [maybeShutdownFinishedFetchers()|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java#L282]
> every time the element queue drains, which reaps any idle fetcher.
> [SingleThreadFetcherManager#addSplits|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SingleThreadFetcherManager.java#L104-L114]
> then allocates a fresh id for the next split, because
> [getRunningFetcher()|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SingleThreadFetcherManager.java#L124-L126]
> finds the map empty.
> Note that the reap path
> [already removes the id from the fetcher
> map|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java#L262]
> -- but the corresponding queue slot is never released:
> {code:java}
> // SplitFetcherManager#createSplitFetcher, L261-L263 -- the shutdown hook
> () -> {
> fetchers.remove(fetcherId);
> fetchersToShutDown.decrementAndGet();
> {code}
> So for any source with continuous split churn, the array becomes a permanent
> tally of every
> fetcher the JVM has created since start.
> h2. Impact
> Two effects, and the second is the more serious one.
> # *Retention.* Each entry holds a {{ConditionAndFlag}} (24 B) and its
> {{AbstractQueuedSynchronizer$ConditionObject}} (24 B), plus a 4 B array slot
> -- 52 B per dead
> fetcher. None of them are reachable by any live thread and none have waiters.
> # *Quadratic copy cost on the hot path.* The array grows *one element at a
> time* via
> {{Arrays.copyOf}}, so total copy work over n fetchers is O(n^2). At the n ~
> 1.2M observed below
> that is roughly 7.3e11 element copies, about 2.9 TB of memcpy over the run.
> Each copy happens
> while holding {{lock}} -- the same lock {{poll()}} acquires -- so the
> consuming {{SourceReader}}
> thread is blocked for its duration. The array also grows toward G1's
> humongous threshold (half
> the region size), at which point every growth becomes a humongous allocation.
> h2. Observed
> Long-running streaming job with high split churn, Flink 2.2.1, JDK 21.0.8,
> G1GC, 4 TaskManagers.
> Three consecutive {{jcmd GC.class_histogram}} samples from one TaskManager.
> All three JFR
> recordings taken alongside them report {{jvmStartTime = 2026-09-08
> 19:10:03.856}} and pid 1, so
> this is a single uninterrupted JVM with no restart between samples.
> {{jdk.CPULoad}} shows the
> JVM was continuously active throughout (jvmUser min 0.57%, mean 3.3-4.3%) --
> it was not idling.
> || sample || JVM uptime || ConditionAndFlag || ConditionObject ||
> ConditionAndFlag[] bytes || live SplitFetcher || live queues ||
> | T0 08:01:29 | 12h 51m | 1,061,289 | 1,062,845 | 4,594,856 | 4 | 4 |
> | T1 10:08:50 | 14h 59m | 1,204,494 | 1,206,050 | 5,208,512 | 4 | 4 |
> | T2 12:52:13 | 17h 42m | 1,388,929 | 1,390,485 | 5,998,680 | 4 | 4 |
> Growth is linear and steady:
> * T0 -> T1: +143,205 over 2h 07m 21s = *67,470/hour*
> * T1 -> T2: +184,435 over 2h 43m 23s = *67,732/hour*
> The decisive detail is that *live {{SplitFetcher}} and
> {{FutureCompletingBlockingQueue}} counts
> stay at 4 across all three samples* -- four live fetchers, four queues,
> unchanged -- while slot
> count grows by 143k per two hours. The {{ConditionObject}} delta is identical
> to the
> {{ConditionAndFlag}} delta on both intervals (+143,205 and +184,435), so
> every retained slot
> pins exactly one {{lock.newCondition()}}. Constant live-fetcher count against
> linear slot growth
> is not explainable by live state.
> At T1 the leaked state accounts for ~63 MB across the four queues, on a 3.6
> GB live set, after
> 15 hours. It does not plateau.
> {noformat}
> taskmanager-1-1, 2026-09-09 08:01:29
> 12: 1062845 25508280
> java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject
> 13: 1061289 25470936
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag
> 18: 4 4594856
> [Lo.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag;
> 1552: 4 288
> o.a.f.c.base.source.reader.fetcher.SplitFetcher
> 2033: 4 160
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue
> Total 75250908 3186293432
> taskmanager-1-1, 2026-09-09 10:08:50
> 12: 1206050 28945200
> java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject
> 13: 1204494 28907856
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag
> 16: 4 5208512
> [Lo.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag;
> 1576: 4 288
> o.a.f.c.base.source.reader.fetcher.SplitFetcher
> 2053: 4 160
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue
> Total 85095888 3599928824
> taskmanager-1-1, 2026-09-09 12:52:13
> 13: 1390485 33371640
> java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject
> 14: 1388929 33334296
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag
> 16: 4 5998680
> [Lo.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue$ConditionAndFlag;
> 1546: 4 288
> o.a.f.c.base.source.reader.fetcher.SplitFetcher
> 2009: 4 160
> o.a.f.c.base.source.reader.synchronization.FutureCompletingBlockingQueue
> Total 97719375 4131002272
> {noformat}
> h2. Possible fixes
> # Release the slot on the path that
> [already removes the fetcher from the
> map|https://github.com/apache/flink/blob/96d7fbe7deaa4ba723d1f08070730e2b4ebe1880/flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/fetcher/SplitFetcherManager.java#L262],
> by also nulling {{putConditionAndFlags[fetcherId]}}. Smallest change, but
> leaves the array length
> at the high-water mark.
> # Recycle fetcher ids through a free list in {{SplitFetcherManager}}, so the
> index stays bounded
> by concurrently-live fetchers. Fixes both the retention and the O(n^2)
> copying.
> # Replace the array with a {{Map<Integer, ConditionAndFlag>}} keyed by
> fetcher id and remove on
> reap. Bounded by live fetchers, no copying at all.
> Happy to put up a PR for whichever direction maintainers prefer.
> h2. Workaround
> Subclassing the fetcher manager to keep a single long-lived {{SplitFetcher}}
> rather than letting
> idle ones be reaped holds the id stable. On our deployment this dropped the
> growth rate from
> ~67,500/hour to ~880/hour, which independently confirms the fetcher-id churn
> as the driver.
> h2. Related
> FLINK-37663 touched the same class (lost wakeup in
> {{FutureCompletingBlockingQueue}}, fixed in
> 2.2.2) but is a separate defect -- a synchronization bug, not unbounded array
> growth. This one is
> not fixed by that change.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)