[
https://issues.apache.org/jira/browse/FLINK-40624?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18116479#comment-18116479
]
Sylwester Lachiewicz commented on FLINK-40624:
----------------------------------------------
[~vbhanuchander] Yes, please take it, I am not working on it. One thing to
check first: [PR #321|https://github.com/apache/flink-connector-kafka/pull/321]
(FLINK-40689) rewrites the commit path in {{KafkaCommitter}} and may already
handle both exceptions, so it is worth agreeing with [~savonitar] which PR
carries this change.
_This comment was created with AI assistance._
> KafkaCommitter and ExactlyOnceKafkaWriter do not handle
> InvalidProducerEpochException alongside ProducerFencedException
> -----------------------------------------------------------------------------------------------------------------------
>
> Key: FLINK-40624
> URL: https://issues.apache.org/jira/browse/FLINK-40624
> Project: Flink
> Issue Type: Bug
> Components: Connectors / Kafka
> Reporter: Sylwester Lachiewicz
> Priority: Minor
>
> A fenced Kafka producer can surface either of two exceptions, depending on
> which request observes the epoch bump:
> * coordinator path (EndTxn, AddPartitions) -> PRODUCER_FENCED ->
> {{ProducerFencedException}}
> * produce path -> INVALID_PRODUCER_EPOCH -> {{InvalidProducerEpochException}}
> In kafka-clients 4.2.0 both extend {{ApplicationRecoverableException}} and
> neither is a subtype of the other. {{InvalidProducerEpochException}} does not
> appear anywhere in flink-connector-kafka, so only the first is handled:
> *
> [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]
> catches {{ProducerFencedException}} and routes it to {{logFencedRequest}}
> plus {{signalFailedWithKnownReason}}. {{InvalidProducerEpochException}}
> misses that branch and every later one, falling through to the terminal
> {{catch (Exception e)}} at
> [KafkaCommitter.java:141|https://github.com/apache/flink-connector-kafka/blob/main/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/internal/KafkaCommitter.java#L141],
> which logs "Transaction ({}) encountered error and data has been potentially
> lost" and calls {{signalFailedWithUnknownReason}}. A fenced producer is then
> reported as an unknown error and the fenced-specific guidance is not emitted.
> *
> [ExactlyOnceKafkaWriter.java:288|https://github.com/apache/flink-connector-kafka/blob/main/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/ExactlyOnceKafkaWriter.java#L288]
> catches {{ProducerFencedException}} to keep the abort best effort, as its
> comment describes. {{InvalidProducerEpochException}} escapes that abort
> instead.
> This is a handling and reporting gap rather than a correctness bug in the
> commit itself; the terminal branch still triggers failover.
> Evidence is the class hierarchy plus the catch chains, and
> [FLINK-40621|https://issues.apache.org/jira/browse/FLINK-40621], where
> {{commitTransaction}} on a fenced producer threw
> {{InvalidProducerEpochException}} with the produce-path message "attempted to
> produce with an old epoch". I have not reproduced the production path.
> Suggested change: handle both exceptions wherever the intent is "the producer
> is fenced", either by catching both or by catching the shared supertype.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)