[
https://issues.apache.org/jira/browse/FLINK-40626?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Sylwester Lachiewicz updated FLINK-40626:
-----------------------------------------
Release Note: With `TransactionNamingStrategy.POOLING`, a recovery now
aborts an open transaction under a precommitted transactional id when the
broker holds a newer epoch for that id. Previously such a transaction stayed
open until `transaction.timeout.ms` whenever the id was never recycled again:
after a downscale, where new transactions only use the first owned subtask id,
or after a change of `transactionalIdPrefix`. `KafkaWriterState` is now
serialized as version 3 (producer id and epoch per precommitted transaction);
state written by earlier versions restores unchanged and keeps the previous
behaviour. Recovery calls `Admin.describeTransactions`, which needs broker 3.0+
and `Describe` on the transactional id. (was: With
`TransactionNamingStrategy.POOLING`, a recovery now aborts an open transaction
under a precommitted transactional id when the broker holds a newer epoch for
that id, instead of leaving it open until `transaction.timeout.ms`.
`KafkaWriterState` is now serialized as version 3 (producer id and epoch per
precommitted transaction); state written by earlier versions restores unchanged
and keeps the previous behaviour. Recovery calls `Admin.describeTransactions`,
which needs broker 3.0+ and `Describe` on the transactional id.)
> KafkaCommitter leaves the newer open transaction under a reused transactional
> id after a fenced commit
> ------------------------------------------------------------------------------------------------------
>
> Key: FLINK-40626
> URL: https://issues.apache.org/jira/browse/FLINK-40626
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Kafka
> Affects Versions: kafka-4.0.1, kafka-5.0.0, kafka-5.1.0
> Reporter: Sylwester Lachiewicz
> Assignee: Sylwester Lachiewicz
> Priority: Major
> Labels: pull-request-available
>
> With {{TransactionNamingStrategy.POOLING}}, a committed transactional id is
> recycled for a later checkpoint;
> {{ProducerPoolImpl.getTransactionalProducer}} calls {{initTransactions()}} on
> reuse, which bumps the epoch. If the job then fails and recovers from the
> earlier checkpoint, the committer re-commits the id at the old epoch and gets
> {{ProducerFencedException}}. That is the documented benign case (the
> transaction was already committed), but the newer transaction the id was
> recycled for is still ONGOING on the broker, and nobody aborts it:
> {{TransactionAbortStrategyImpl.LISTING}} skips the id because it is in the
> recovered precommitted set
> ([TransactionAbortStrategyImpl.java:145|https://github.com/apache/flink-connector-kafka/blob/main/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/TransactionAbortStrategyImpl.java#L145]),
> and the committer drops the request
> ([KafkaCommitter.java:109|https://github.com/apache/flink-connector-kafka/blob/main/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/KafkaCommitter.java#L109]).
> The open transaction pins the last stable offset until
> {{transaction.timeout.ms}} (1 h by default), so {{read_committed}} consumers
> see nothing the recovered job writes. With an unchanged prefix the writer
> eventually reuses the id and the epoch bump aborts it (a stall of about two
> checkpoint intervals). With a changed {{transactionalIdPrefix}} (supported;
> the writer logs it and adds the old prefix to the abort list) the id is never
> touched again and the stall lasts until the timeout.
> Observed in CI as
> [FLINK-40585|https://issues.apache.org/jira/browse/FLINK-40585] (run
> 30181683657, {{rescaleListing[3->5]}}): job 3 lists {{secondPrefix-0-0}} as
> open, skips it as precommitted, the committer's re-commit at epoch 11 is
> fenced, and the drain stops after job 2's records.
> Proposed fix: on a fenced commit with reused ids, the committer bumps the
> epoch of that id ({{initTransactions}} on a short-lived producer) before
> dropping the request, which aborts the orphaned transaction.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)