Alexandru created FLINK-40657:
---------------------------------

             Summary:  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.1.3, 2.2.1, 2.3.0, 2.0.2, 1.20.6
            Reporter: Alexandru


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

Reply via email to