Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

...

Exception/Error Names

Current handling

New Handling


Producer API

Transaction API

Producer API

Transaction API

IllegalStateException

Abortable

Sometimes Fatal depending on whether application or Sender caused issue. See: kafka: KAFKA-14831: Illegal state errors should be fatal in transactional producerCLOSED

Application Recoverable (probably not expected)

Application Recoverable

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

Abortable

Application Recoverable

Application Recoverable

FencedInstanceIdException

CommitFailedException

UnknownMemberIdException

IllegalGenerationExceiption

N/A

Abortable (TxnOffsetCommit Only)

N/A

Application Recoverable

KafkaException

Abortable (default seems to be abortable)

Fatal in most cases, but abortable when there are partition errors

Application Recoverable (not expected)

Application Recoverable (not expected)

RuntimeException

Abortable (default seems to be abortable)

Fatal – only thrown as this generic type when correlation ID is wrong. This should be updated as KIP-691 suggests

Application Recoverable (not expected)

Application Recoverable (not expected)

CorrelationIdMismatchException



Application Recoverable

Application Recoverable

...

InvalidConfigurationTransactionException

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

...

Code Block
// Producer Abortable Transaction
public class AbortableTransactionException extends ApiException {
    ...
}

//Producer Retriable 
public class ProducerRetriableTransactionException extends ApiException {
 	...
}

//Producer Refresh and Retriable
public class ProducerRefreshRetriableTransactionException extends ApiException {
	...
}

//Application-Recoverable
public class ApplicationRecoverableTransactionException extends ApiException {
    ...
}

// Invalid-Configuration
public class InvalidConfiguationTransactionExceptionInvalidConfigurationTransactionException extends ApiException {
 	...
}

...

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) {
                            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 (AbortableTransactionException e) {
                    // Abortable Exception: Handle Kafka exception by aborting transaction. producer.abortTransaction() should not throw abortable exception.
                    producer.abortTransaction();
                    resetToLastCommittedPositions(consumer);
                }
            } catch (InvalidConfiguationTransactionExceptionInvalidConfigurationTransactionException e) {
                //  Fatal Error: The error is bubbled up to the application layer. The application can decide what to do
                closeAll();
                throw e;
            } catch (KafkaException | ApplicationRecoverableTransactionException e) {
                // Application Recoverable: The application must restart
                closeAll();
                initializeApplication();
            }
        }

    }

...