DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
Retry the operation without bothering the application (retriable)
Share with the application that an error has occurred, but the client can recover if the transaction is aborted by the client and under the covers, the epoch is bumped (abortable)
Share with the application that an error has occurred, and there is no way to recover except shut down the client (fatal)
...
We are proposing to group exceptions into four types:
Producer-Retriable: mostly maps onto current retriable errors. Some of these errors may indicate that a metadata update at the producer level is required. Because of that we introduce a few two subclasses to the retriable error
retry only (just send the request again – after some period of backoff)
refresh metadata and retry (request metadata and maybe modify the request before resending)
- Producer-Abortable: mostly maps onto abortable errors. The error is bubbled to the application layer, and it can choose to roll back any state from the ongoing transaction and abort the transaction. The producer does not need to restart, as after aborting, it can be confident that the state was as it was before the transaction started. The application can also close the producer and react as it does for application-recoverable cases, but it doesn’t need to.
Application-Recoverable: maps on to some fatal errors. The error is bubbled to the application layer, and it may need to do a bit more to roll back and clean up state. The producer can not simply abort and know the current state of the partitions. Applications can handle this in different ways – Kafka Streams may rebalance a task and/or close the task and restart it. Another application may read from their own checkpoint to continue. In any case, the producer must restart and will be unusable after encountering this error.
Invalid-Configuration: maps to some fatal errors. The error is bubbled up to the application layer. The application can decide what to do. The producer doesn’t need to restart, but the application may chose to close it.
For Producer-Retriable errors, the producer handles retries internally, keeping the failure details hidden from the application. Conversely, other types of exceptions will be surfaced to the application code for handling.
Each error code always represent the same class and rarely rely on client state to determine how to handle. Additionally, while it is good to have a mapping, It is also useful to have a general strategy – ie a typical unknown error (not specified by the client to have a type) should probably be application recoverable.
Exception Table
NOTE: For Retriable errors, the producer handles retries internally, keeping the failure details hidden from the application. Conversely, other types of exceptions will be surfaced to the application code for handling.
Exception Table
Below table contains Below table contains list of exceptions with current and New handing. All exceptions are grouped under a header which represent Exception Name:
...
Cross → Exception handling doesn't match as mentioned in KIP-691
...
Exception/Error Names | Current handling | New Handling | ||
|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | |
CorruptRecordException NotEnoughReplicasAfterAppendException NotEnoughReplicasException TimeoutException | Retriable if the error is retriable. Otherwise abortable | Retriable if the error is retriable otherwise abortable | Producer Retriable | Producer Retriable |
ConcurrentTransactionsException | Retriable | Retriable | Producer Retriable (KIP-890 may see this) | Producer Retriable |
CoordinatorLoadInProgressException | Retriable | Retriable | Producer Retriable | Producer Retriable |
OutOfOrderSequenceException | Sometimes retriable or abortable | N/A | Producer Retriable | N/A |
Exception/Error Names | Current handling | New Handling | ||
|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | |
UnknownTopicOrPartitionException NotLeaderOrFollowerException | Retriable if the error is retriable. Otherwise abortable | Retriable if the error is retriable otherwise abortable | Refresh + Retriable | Refresh + Retriable |
NotCoordinatorException CoordinatorNotAvailableException | Retriable | Retriable | Refresh + Retriable | Refresh + Retriable |
...
Exception/Error Names | Current handling | New Handling | ||
|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | |
TransactionAbortableException | Producer Abortable | Producer Abortable | Producer Abortable | Producer Abortable |
...
ApplicationRecoverableException
Exception/Error Names | Current handling | New Handling | |||||||
|---|---|---|---|---|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | ||||||
InvalidProducerEpochException | Abortable | Fatal | Application Recoverable | Application Recoverable | |||||
ProducerFencedException | Fatal | Fatal | Application Recoverable | Application Recoverable | |||||
InvalidPidMappingException | Sometimes retriable (no longer returned) or abortable | Abortable | Application Recoverable | Application Recoverable | UnknownProducerIdException | Sometimes retriable (no longer returned) or abortable | AbortableApplication Recoverable | Application Recoverable | |
FencedInstanceIdException CommitFailedException UnknownMemberIdException IllegalGenerationExceiption | N/A | Abortable (TxnOffsetCommit Only) | N/A | Application Recoverable | |||||
CorrelationIdMismatchException | Application Recoverable | Application Recoverable | |||||||
...
Exception/Error Names | Current handling | Expected Handling | ||
|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | |
InvalidTxnStateException | Abortable | Fatal | Producer Abortable (KIP-890 relies on this – note this is the only one that differs with Produce API) | Application Recoverable |
...
InvalidConfigurationException
Exception/Error Names | Current handling | Expected Handling | ||
|---|---|---|---|---|
Producer API | Transaction API | Producer API | Transaction API | |
AuthenticationException | Abortable | Fatal | Invalid Configuration | Invalid Configuration (not expected) |
ClusterAuthorizationException TransactionalIdAuthorizationException UnsupportedVersionException UnsupportedForMessageFormatException | Fatal, except UnsupportedForMessageFormatException which is Abortable | Cluster/Transaction Auth → abortable on InitProducerId, other errors abortable TransactionAuth → fatal, others abortable for AddPartitions, Find Coordinator, EndTxn, AddOffsets Transaction Auth, UnssupportedForMessageFormat → fatal for OffsetCommit | Invalid Configuration | Invalid Configuration |
InvalidRecordException InvalidRequiredAcksException RecordBatchTooLargeException InvalidTopicException | Retriable if the error is retriable. Otherwise abortable | Retriable if the error is retriable otherwise abortable | Invalid Configuration | Invalid Configuration |
TopicAuthorizationException GroupAuthorizationException | Abortable | Abortable | Invalid Configuration | Invalid Configuration |
Public Interfaces
We will add new exception types, as listed in the table below, that extend the existing exceptions.have below exception classes available in current Kafka code
| Code Block |
|---|
// Producer Abortable Transaction public class AbortableTransactionExceptionTransactionAbortableException extends ApiException { ... } //Producer Retriable public abstract class ProducerRetriableTransactionExceptionRetriableException extends ApiException { ... } //Producer Refresh and Retriable Invalid-Configuration public class ProducerRefreshRetriableTransactionExceptionInvalidConfigurationException extends ApiException { ... } //Application-Recoverable public class ApplicationRecoverableTransactionException extends ApiException { |
We will add new exception types, as listed in the table below, that extend the existing exceptions.
| Code Block |
|---|
//Producer Refresh and Retriable public abstract class RefreshRetriableException extends RetriableException { ... } // Invalid-ConfigurationApplication-Recoverable new public abstract class InvalidConfigurationTransactionExceptionApplicationRecoverableException extends ApiException { ... } |
Client side code example
| Code Block |
|---|
public class TransactionalClientDemo {
private static final String CONSUMER_GROUP_ID = "my-group-id";
private static final String OUTPUT_TOPIC = "output";
private static final String INPUT_TOPIC = "input";
private static KafkaConsumer<String, String> consumer;
private static KafkaProducer<String, String> producer;
public static void main(String[] args) {
initializeApplication();
boolean isRunning = true;
// Continuously poll for records
while (isRunning) {
try {
try {
// Poll records from Kafka for a timeout of 60 seconds
ConsumerRecords<String, String> records = consumer.poll(ofSeconds(60));
// Process records to generate word count map
Map<String, Integer> wordCountMap = new HashMap<>();
for (ConsumerRecord<String, String> record : records) {
String[] words = record.value().split(" ");
for (String word : words) {
UnknownProducerIdException wordCountMap.merge(word, 1, Integer::sum);
}
}
// Begin transaction
producer.beginTransaction();
// Produce word count results to output topic
wordCountMap.forEach((key, value) ->
producer.send(new ProducerRecord<>(OUTPUT_TOPIC, key, value.toString())));
// Determine offsets to commit
Map<TopicPartition, OffsetAndMetadata> offsetsToCommit = new HashMap<>();
for (TopicPartition partition : records.partitions()) {
List<ConsumerRecord<String, String>> partitionedRecords = records.records(partition);
long offset = partitionedRecords.get(partitionedRecords.size() - 1).offset();
offsetsToCommit.put(partition, new OffsetAndMetadata(offset + 1));
}
// Send offsets to transaction for atomic commit
producer.sendOffsetsToTransaction(offsetsToCommit, CONSUMER_GROUP_ID);
// Commit transaction
producer.commitTransaction();
} catch (AbortableTransactionExceptionTransactionAbortableException e) {
// Abortable Exception: Handle Kafka exception by aborting transaction. producer.abortTransaction() should not throw abortable exception.
producer.abortTransaction();
resetToLastCommittedPositions(consumer);
}
} catch (InvalidConfigurationTransactionExceptionInvalidConfigurationException e) {
// Fatal Error: The error is bubbled up to the application layer. The application can decide what to do
closeAll();
throw e;
} catch (KafkaException | ApplicationRecoverableTransactionExceptionApplicationRecoverableException e) {
// Application Recoverable: The application must restart
closeAll();
initializeApplication();
}
}
} |
...
For existing errors, if the application code already handles the error, it will continue to do so. However, the new TransactionAbortableException introduced in KIP-890 extends KafkaException, which is treated as fatal by applications. Therefore, we expect TransactionAbortableException to also be treated as fatal by older code TransactionAbortableException to also be treated as fatal by older code.
Currently, the transactional producer returns retriable exception types, such as TimeoutException , which poses a risk of duplicates in Kafka. In this KIP, we will update the transactional producer to return TransactionAbortableException instead of TimeoutException . Older clients that are using the transactional producer and handling TimeoutException by retrying the produce operation can update to handle TransactionAbortableException .
For applications using newer code
Applications using newer code, as mentioned in Clientsidecodeexample, have the flexibility to adjust their exception handling logic based on the specific type of exception encountered, enabling them to take appropriate actions.
Handling for OutOfOrderSequenceException requires additional considerations which is out of scope for this KIP. OutOfOrderSequenceException andUnknownProducerIdException(which extends OutOfOrderSequenceException) will not be categorised under any of the defined exception groups in this KIP. This will not impact applications using newer code.
Test Plan
We will add additional integration and unit tests to check if client acts as expected.
...