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).
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:
This KIP introduces an interface that can be implemented by the user to handle the exceptions UnknownTopicOrPartitionException and RecordTooLargeException.
Question: Why do we need an interface for handling the exceptions? Could we have a couple of simple producer configuration options for those two exceptions?
Answer: We aim at giving the user the flexibility of an interface. For example, facing UnknownTopicOrPartitionException, the user may want to raise an error for some topics but retry it for other topics. Having a configuration option with a fixed set of possibilities does not serve the user's needs.
We introduce the ProducerExceptionHandler interface, that can be implemented by the user to manage the exception in the desired manner.
To configure their own handler, the user must implement the above introduced interface and add the class name in producer configuration with the key: custom.exception.handler.
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
*/
Response handle(final ProducerRecord<byte[], byte[]> record,
final Exception exception);
enum Response {
/* 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;
}
}
}
|
.
.
.
public static final String CUSTOM_EXCEPTION_HANDLER_CLASS_CONFIG = "custom.exception.handler";
private static final String CUSTOM_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(CUSTOM_EXCEPTION_HANDLER_CLASS_CONFIG,
Type.CLASS,
null,
Importance.MEDIUM,
CUSTOM_EXCEPTION_HANDLER_CLASS_DOC);
} |
The custom handler will only affect the exceptions thrown from the producer. This KIP, very specifically means to provide a possibility for users to manage a limited number of exceptions (RecordTooLargeException and UnknownTopicOrPartitionException so far) thrown from the producer send() method. Of course the same exceptions may originate from different components of Apache Kafka which are not the focus of this KIP.
Notes on RecordTooLargeException based on the codebase we have today:
max.request.size" and "buffer.memory" to big numbers in producer config), but it is too large for the broker (The default message size of broker is 1 MB). In such case, the broker throws RecordTooLargeException during commitTransaction(). This scenario is not the focus of this KIP.
Compatibility, Deprecation, and Migration Plan
Changed behaviour: The default behaviour stays as it is, but the user can change the behaviour by implementing the handle() function.
Some unit and integration tests will be implemented to ensure that
No unit tests are needed to ensure the backward compatibility. Passing the current unit tests is an enough indicator.
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.