DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
| Code Block | ||
|---|---|---|
| ||
@Override
public ProcessingHandlerResponse handleError(final ErrorHandlerContext context, final Record<?, ?> record, final Exception exception) {
List<ProducerRecord<byte[], byte[]>> records = Collections.singletonList(new ProducerRecord<>("app-dlq", "Hello".getBytes(StandardCharsets.UTF_8), "World".getBytes(StandardCharsets.UTF_8)));
return ProcessingExceptionResponse.continueProcessing(records); } |
...
| Code Block | ||
|---|---|---|
| ||
public interface ProductionExceptionHandler extends Configurable {
... /**
@Deprecated
default ProductionExceptionHandlerResponse handle(final ErrorHandlerContext context,
final ProducerRecord<byte[], byte[]> record,
final Exception exception) {
throw new UnsupportedOperationException();
}
/**
* Inspect a record that we attempted to produce, and the exception that resulted
* from attempting to produce it and determine to continue or stop processing.
*
* @param context
* The error handler context metadata.
* @param record
* The record that failed to produce.
* @param exception
* The exception that occurred during production.
*
* @return a {@link ProductionExceptionResponse} object
*/
default ProductionExceptionResponse handleError(final ErrorHandlerContext context,
final ProducerRecord<byte[], byte[]> record,
final Exception exception) {
final ProductionExceptionHandlerResponse response = handle(context, record, exception);
if (ProductionExceptionHandler.ProductionExceptionHandlerResponse.FAIL == response) {
return ProductionExceptionResponse.failProcessing();
} else if (ProductionExceptionHandler.ProductionExceptionHandlerResponse.RETRY == response) {
return ProductionExceptionResponse.retryProcessing();
}
return ProductionExceptionResponse.continueProcessing();
}
@Deprecated
default ProductionExceptionHandlerResponse handleSerializationException(final ProducerRecord record,
final Exception exception) {
return ProductionExceptionHandler.ProductionExceptionHandlerResponse.FAIL;
}
/**
* Handles serialization exception and determine if the process should continue. The default implementation is to
* fail the process.
*
* @param context
* The error handler context metadata.
* @param record
* The record that failed to serialize.
* @param exception
* The exception that occurred during serialization.
* @param origin
* The origin of the serialization exception.
*
* @return a {@link ProductionExceptionResponse} object
*/
default ProductionExceptionResponse handleSerializationError(final ErrorHandlerContext context,
final ProducerRecord record,
final Exception exception,
final SerializationExceptionOrigin origin) {
final ProductionExceptionHandlerResponse response = handleSerializationException(context, record, exception, origin);
if (ProductionExceptionHandler.ProductionExceptionHandlerResponse.FAIL == response) {
return ProductionExceptionResponse.failProcessing();
} else if (ProductionExceptionHandler.ProductionExceptionHandlerResponse.RETRY == response) {
return ProductionExceptionResponse.retryProcessing();
}
return ProductionExceptionResponse.continueProcessing();
}
...
/**
* Represents the result of handling a production exception.
* <p>
* The {@code Response} class encapsulates a {@link ProductionExceptionHandlerResponse},
* indicating whether processing should continue or fail, along with an optional list of
* {@link ProducerRecord} instances to be sent to a dead letter queue.
* </p>
*/
class ProductionExceptionResponse {
private ProductionExceptionHandlerResponse productionExceptionHandlerResponse;
private List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords;
/**
* Constructs a new {@code ProductionExceptionResponse} object.
*
* @param productionExceptionHandlerResponse the response indicating whether processing should continue or fail;
* must not be {@code null}.
* @param deadLetterQueueRecords the list of records to be sent to the dead letter queue; may be {@code null}.
*/
private ProductionExceptionResponse(final ProductionExceptionHandlerResponse productionExceptionHandlerResponse,
final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
this.productionExceptionHandlerResponse = productionExceptionHandlerResponse;
this.deadLetterQueueRecords = deadLetterQueueRecords;
}
/**
* Creates a {@code ProductionExceptionResponse} indicating that processing should fail.
*
* @param deadLetterQueueRecords the list of records to be sent to the dead letter queue; may be {@code null}.
* @return a {@code ProductionExceptionResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#FAIL} status.
*/
public static ProductionExceptionResponse failProcessing(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
return new ProductionExceptionResponse(ProductionExceptionHandlerResponse.FAIL, deadLetterQueueRecords);
}
/**
* Creates a {@code ProductionExceptionResponse} indicating that processing should fail.
*
* @return a {@code ProductionExceptionResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#FAIL} status.
*/
public static ProductionExceptionResponse failProcessing() {
return new ProductionExceptionResponse(ProductionExceptionHandlerResponse.FAIL, failProcessing(Collections.emptyList());
}
/**
* Creates a {@code ProductionExceptionResponse} indicating that processing should continue.
*
* @param deadLetterQueueRecords the list of records to be sent to the dead letter queue; may be {@code null}.
* @return a {@code ProductionExceptionResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
*/
public static ProductionExceptionResponse continueProcessing(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
return new ProductionExceptionResponse(ProductionExceptionHandlerResponse.CONTINUE, deadLetterQueueRecords);
}
/**
* Creates a {@code ProductionExceptionResponse} indicating that processing should continue.
*
* @return a {@code ProductionExceptionResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
*/
public static ProductionExceptionResponse continueProcessing() {
return new ProductionExceptionResponse(ProductionExceptionHandlerResponse.CONTINUE, continueProcessing(Collections.emptyList());
}
/**
* Creates a {@code ProductionExceptionResponse} indicating that processing should retry.
*
* @return a {@code ProductionExceptionResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
*/
public static ProductionExceptionResponse retryProcessing() {
return new ProductionExceptionResponse(ProductionExceptionHandlerResponse.RETRY, Collections.emptyList());
}
/**
* Retrieves the production exception handler response.
*
* @return the {@link ProductionExceptionHandlerResponse} indicating whether processing should continue or fail.
*/
public ProductionExceptionHandlerResponse response() {
return productionExceptionHandlerResponse;
}
/**
* Retrieves an unmodifiable list of records to be sent to the dead letter queue.
* <p>
* If the list is {@code null}, an empty list is returned.
* </p>
*
* @return an unmodifiable list of {@link ProducerRecord} instances
* for the dead letter queue, or an empty list if no records are available.
*/
public List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords() {
if (deadLetterQueueRecords == null) {
return Collections.emptyList();
}
return Collections.unmodifiableList(deadLetterQueueRecords);
}
}
...
}
|
...