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

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

> KafkaSinkITCase.rescaleListing fails intermittently after scale-up-then-down: 
> restored transactions fenced on commit, only a subset of records readable
> -------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40585
>                 URL: https://issues.apache.org/jira/browse/FLINK-40585
>             Project: Flink
>          Issue Type: Bug
>          Components: Connectors / Kafka
>    Affects Versions: kafka-4.0.1, kafka-5.1.0
>            Reporter: Martijn Visser
>            Priority: Major
>              Labels: pull-request-available
>
> - run 30181683657, job 89739208487 (2026-07-26, main, Flink 2.2-SNAPSHOT, JDK 
> 21), rescaleListing[4] (3 -> 5 -> 3): actual [334L, 335L, 336L, 667L, 668L, 
> 669L, 0L, 1L, 2L, 3L, 4L, 337L, 338L, 670L, 671L], expected [0L .. 999L].
>  - run 31916828021, job 95089941850 (2026-08-16, v4.0, Flink 2.0.2, JDK 21), 
> rescaleListing[3] (2 -> 5 -> 2): 966 of 1000 records missing.
> The test uses TransactionNamingStrategy.POOLING; job 1 (firstPrefix) fails in 
> the async checkpoint phase, job 2 (secondPrefix) restores at a different 
> parallelism and fails the same way, job 3 (thirdPrefix) restores from job 2's 
> checkpoint at the original parallelism.
> What the job-3 log shows (89739208487, 01:07:53):
> .648 committers resume secondPrefix-0-0 / -2-0 / -4-0 with the restored 
> producerId and epoch 0
> .756 KafkaCommitter WARN "Unable to commit transaction ... because its 
> producer is already fenced" for all three (ProducerFencedException), ~100 ms 
> later
> .770 writers start "Aborting lingering transactions with prefixes 
> [thirdPrefix, secondPrefix]
> using LISTING"; the listing finds only secondPrefix-x-1 open, i.e. the -x-0 
> transactions
> were no longer open on the broker at that point, and the precommitted ids 
> were correctly
> skipped by TransactionAbortStrategyImpl
> 01:08:09 the retry creates fresh producers (initTransactions, epoch 2) and 
> the commit fails
> again with ERROR.
> So the epoch of the precommitted transactions had already moved past 0 before 
> job 3 started. Two readings, and the log cannot separate them:
> (a) benign: job 2 had already committed them on checkpoint completion 
> (RecordFetcher saw the notification), the POOLING recycle bumped the epoch, 
> and the fenced retry is the documented "expected in some cases" path. The 
> records are then in the topic, and the missing data is a read_committed 
> artefact: KafkaUtil.drainAllRecordsFromTopic snapshots endOffsets once (the 
> LSO), and job 3 fails with "Planned exception" while its thirdPrefix-x-1 
> transactions stay open (the writer deliberately does not abort on close, 
> ExactlyOnceKafkaWriter ~line 281), pinning the LSO.
> (b) real: the transactions were never committed and the committer dropped 
> them via
> signalFailedWithKnownReason, which the warning text itself says signals data 
> loss.
> Reading (a) predicts that draining with read_uncommitted, or aborting the 
> leftover thirdPrefix transactions before draining, would show all 1000 
> records; reading (b) predicts they are absent. That single check decides it 
> and should be the first step. Until then treat as Major, not Critical; if (b) 
> holds, raise to Blocker and it is a FLIP-511 recovery bug.



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

Reply via email to