You are viewing an old version of this page. View the current version.

Compare with Current View Page History

« Previous Version 15 Next »

Status

Current state: Under discussion 

Discussion thread: here

JIRA: here 

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

Motivation

Kafka Producer supports following transactional APIs

  1. initTransactions for transactional producer identity initialization

  2. beginTransaction to start a new transaction

  3. sendOffsetsToTransaction to commit consumer offsets advanced within the current transaction

  4. commitTransaction commit the ongoing transaction

  5. abortTransaction abort the ongoing transaction

In addition to above transactional APIs, idempotent producer sends records to broker using send API


Apache Kafka supports a variety of client SDKs for different programming languages:

  1. Java: The official Java client library supports the producer, consumer, Streams, and Connect APIs.

  2. librdkafka and derived clients:

    1. C/C++: A C/C++ client library supporting the Producer and Consumer APIs.

    2. Python: A Python client library supporting the Producer and Consumer APIs.

    3. Go: A Go client library supporting the Producer and Consumer APIs.

    4. .NET: A .NET client library supporting the Producer and Consumer APIs.


We require consistent error handling across all clients SDKs and APIs. Functionally the client can handle error by:

  • 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)


In this KIP, we want to accomplish a few goals:

We have attempted error handing in accepted KIP: KIP-691: Enhance Transactional Producer Exception Handling - Apache Kafka - Apache Software Foundation which we can perhaps leverage some of the work here.

Public Interfaces

We would new error types which will be extended by existing exceptions mentioned in below table.

// Producer-Recoverable
public class AbortableTransactionException extends ApiException {
    public AbortableTransactionException(String message) {
        super(message);
    }
    ...
}

//Producer-Retriable 
public class ProducerRetriableTransactionException extends ApiException {
    public ProducerRetriableTransactionException(String message) {
        super(message);
    }
 	...
}

//Producer-Retriable
public class ProducerImmediateRetriableTransactionException extends ApiException {
    public ProducerImmediateRetriableTransactionException(String message) {
        super(message);
    }
	...
}

//Application-Recoverable
public class ApplicationRecoverableTransactionException extends ApiException {
    public ApplicationRecoverableTransactionException(String message) {
        super(message);
    }
    ...
}

// Invalid-Configuration
public class InvalidConfiguationTransactionException extends ApiException {
    public InvalidConfiguationTransactionException(String message) {
        super(message);
    }
 	...
}

// Extending exception types example
public class InvalidProducerEpochException extends AbortableTransactionException {
    private static final long serialVersionUID = 1L;
    public InvalidProducerEpochException(String message) {
        super(message);
    }
}


Client side code example

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 (InvalidConfiguationTransactionException 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();
            }
        }

    }

Full example can be accessed at: https://github.com/apache/kafka/pull/15913/files 


Proposed Changes

We are proposing to group exceptions into four types:

  1. Producer-Retriable: mostly maps onto current retriable errors. The producer can handle the retry on its own and the failure is invisible to the application. Some of these errors may indicate that a metadata update at the producer level is required. Because of that we introduce a few subclasses to the retriable error

    1. retry only (just send the request again – after some period of backoff)

    2. refresh metadata and retry (request metadata and maybe modify the request before resending)

  2. Producer-Recoverable: 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.

  3. 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 – 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.

  4. 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.

Each error code always represent the same class and rarely rely on client state to determine how to handle.

While it is good to have a mapping, I think 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.


Below table contains list of exceptions with current and expected handing:

Yellow → Needs modification

Green → New exceptions added as part of this KIP

Check Exception handling match as mentioned in KIP-691

Cross Exception handling doesn't match as mentioned in KIP-691


Exception Table 

Exception/Error Names

Thrown scenarios during transaction

Current handling

Expected Handling



Producer API

Transaction API

Producer API

Transaction API

ErrorRetry


Producer Retriable

Producer Retriable

Producer Retriable

Producer Retriable

ErrorRefreshMetadataAndRetry


Refresh + Retriable

Refresh + Retriable

Refresh + Retriable

Refresh + Retriable

ErrorFenceClient


Application Recoverable

Application Recoverable

Application Recoverable

Application Recoverable

ErrorBadConfig


Invalid Configuration

Invalid Configuration

Invalid Configuration

Invalid Configuration

TransactionAbortableException
** note – Added via KIP-890


Producer Recoverable

Producer Recoverable

Producer Recoverable

Producer Recoverable

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

Au thenticationException


Abortable

Fatal

Invalid Configuration

Invalid Configuration (not expected)

InvalidPidMappingException

  • If the TransactionalId does not exist in the TransactionalId mapping or if the mapped PID is different from that in the request, reply with InvalidPidMapping; otherwise proceed to next step.

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

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

CorruptRecordException

NotEnoughReplicasAfterAppendException

NotEnoughReplicasException

TimeoutException


Retriable if the error is retriable. Otherwise abortable

Retriable if the error is retriable otherwise abortable

Producer Retriable

Producer Retriable

UnknownTopicOrPartitionException

NotLeaderOrFollowerException


Retriable if the error is retriable. Otherwise abortable

Retriable if the error is retriable otherwise abortable

Refresh + Retriable

Refresh + Retriable

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

FencedInstanceIdException

CommitFailedException

UnknownMemberIdException

IllegalGenerationExceiption


N/A

Abortable (TxnOffsetCommit Only)

N/A

Application Recoverable

InvalidProducerEpochException

  • If the PID’s epoch number is different from the current TransactionalId PID mapping, reply with the InvalidProducerEpoch error code; otherwise proceed to next step.

Abortable

Fatal

Application Recoverable

Application Recoverable

ProducerFencedException


Fatal

Fatal

Application Recoverable

Application Recoverable

OutOfOrderSequenceException


Sometimes retriable or abortable

N/A

Producer Retriable

N/A

InvalidTxnStateException

  • When client makes a request which involves change of Txn state. However expected state transition is invalid.
    e.g COMPLETE_ABORT to COMMIT

Abortable

Fatal

Producer Recoverable (KIP-890 relies on this – note this is the only one that differs with Produce API)

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)

ConcurrentTransactionsException


Retriable

Retriable

Producer Retriable (KIP-890 may see this)

Producer Retriable

CoordinatorLoadInProgressException


Retriable

Retriable

Producer Retriable

Producer Retriable

NotCoordinatorException

CoordinatorNotAvailableException

  • Check whether there is a previous entry with the same TransactionalId and a higher epoch. If so, throw an exception. In particular, this indicates the log is corrupt. All future transactional RPCs to this coordintaor will result in a `NotCoordinatorForTransactionalId` error code, and this partition of the log will be effectively disabled.

  • Check if it is the assigned transaction coordinator for the TransactionalId, if not reply with the NotCoordinatorForTransactionalId error code.

Retriable

Retriable

Refresh + Retriable

Refresh + Retriable

CorrelationIdMismatchException




Application Recoverable

Application Recoverable

Compatibility, Deprecation, and Migration Plan

This KIP involves client side changes which only affects the resiliency of new Producer client and Streams. Old clients will continue to be able to use their error handling. Based on type of exception thrown, user needs to change their exception catching logic to take actions against their exception handling.

Test Plan

We will add additional integration and unit tests to check if client acts as expected.

Rejected Alternatives

We discussed another approach where we encode the handling into the response as a separate field from the error code or something similar. Produce responses have space for a code and for a message, so something similar could be done to distinguish the cause/meaning of the error from the action that should be taken. The encoding of the error could be included in the definition of the response spec for other client compatibility. We rejected this approach as it requires a larger overhaul on the APIs as well as the clients.


  • No labels