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

Compare with Current View Page History

« Previous Version 11 Next »

Status

Current state: Under Discussion

Discussion thread: here

JIRA: KAFKA-15309 
Other related tickets: KAFKA-9279, KAFKA-15259

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

Motivation

We believe that the user should be able to develop custom exception handlers for managing producer exceptions. On the other hand, this will be an expert-level API, and using that may result in strange behaviour in the system, making it hard to find the root cause. Therefore, the custom handler is currently limited to handling RecordTooLargeException and UnknownTopicOrPartitionException. The motivation for this KIP is derived from the following use cases:

  • In transactions, the producer collects multiple records in batches. Then a RecordTooLargeException related to a single record leads to failing the entire batch. A custom exception handler in this case may decide on dropping the record and continuing the processing.
  • When a user tries to write into a non-existing topic, it returns a retryable error code; with infinite retries, the producer would hang retrying forever. A custom handler in this case could help to break the infinite retry loop.

Public Interfaces

We introduce the ProducerExceptionHandler interface, which can be implemented by the user to manage the exception in the desired manner.

ProducerExceptionHandler
package org.apache.kafka.common.errors;

import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.Configurable;
import org.apache.kafka.common.annotation.InterfaceStability;

/**
 * Interface that specifies how an exception should be handled.
 */
@InterfaceStability.Evolving
public interface ProducerExceptionHandler extends Configurable {

    /**
     * Determine whether to stop processing, keep retrying internally, or swallow the error by dropping the record.
     *
     * @param record The record that failed to produce
     * @param exception The exception that occurred during production
     */
    ProducerExceptionHandlerResponse handle(final ProducerRecord<byte[], byte[]> record,
                                            final Exception exception);

    enum ProducerExceptionHandlerResponse {
        /* stop processing: fail */
        FAIL(0, "FAIL"),
        /* continue: keep retrying */
        RETRY(1, "RETRY"),
        /* continue: swallow the error */
        SWALLOW(2, "SWALLOW");

        /**
         * an english description of the api--this is for debugging and can change
         */
        public final String name;

        /**
         * the permanent and immutable id of an API--this can't change ever
         */
        public final int id;

        ProducerExceptionHandlerResponse(final int id,
                                         final String name) {
            this.id = id;
            this.name = name;
        }
    }
}
 


ProducerConfig
public static final String PRODUCER_EXCEPTION_HANDLER_CLASS_CONFIG = "producer.exception.handler";
private static final String PRODUCER_EXCEPTION_HANDLER_CLASS_DOC = "Exception handling class that implements the <code>org.apache.kafka.common.errors.ProducerExceptionHandler</code> interface.";.
.
.
 static {
CONFIG = new ConfigDef().define(
.....
.
.
.
                                              .define(PRODUCER_EXCEPTION_HANDLER_CLASS_CONFIG,
                                                           Type.CLASS,
                                                           null,
                                                           Importance.MEDIUM,
                                                           PRODUCER_EXCEPTION_HANDLER_CLASS_DOC);}


Compatibility, Deprecation, and Migration Plan

To configure own handler, the user must implement the above introduced interface and add the class name in producer configuration with the key: producer.exception.handler.

Changed behaviour: The default behaviour stays as it is, but the user can change the behaviour by implementing the handle() function.

Test Plan

Some unit and integration tests will be implemented to ensure that

  • exceptions are caught by the ProducerExceptionHandler and the provided implementations.
  • the right exceptions are caught.

No unit tests are needed to ensure the backward compatibility. Passing the current unit tests is an enough indicator.

Rejected Alternatives

Using one or more producer configs instead of having a pluggable interface: misusing produce configs has the same drawbacks of misusing the interface while the interface solution provides a handler with the advantage of full flexibility. Later, further KIPs can be proposed to cover more exceptions or more actions for handing. 



  • No labels