[
https://issues.apache.org/jira/browse/FLINK-40626?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Martijn Visser reassigned FLINK-40626:
--------------------------------------
Assignee: Sylwester Lachiewicz
> 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)