Versions Compared

Key

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

...

ProductionExceptionHandler.java

Changes:

  • Adding a ProductionExceptionResponse Response nested class,that contains a ProductionExceptionHandlerResponse indicating whether to continue processing, fail, or retry and the list of records to be sent to the dead letter queue topic
  • Deprecating the handler() and handlerSerializationException() methods and adding two new methods: handlerError() and handleSerializationException(). The return type is the ProductionExceptionResponse
  • As the ErrorHandlerContext does not provide the sourceKey/Value in the handle method to limit the memory impact, the Dead Letter Queue record would only contains metadata. handleSerializationException is not impacted.
Code Block
languagejava
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) {                 
        return new ProductionExceptionResponse(handle(context, record, exception), Collections.emptyList());
     }

    @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) {         
      return new ProductionExceptionResponse(handleSerializationException(context, record, exception, origin), Collections.emptyList());     
    }

    ...

     /**
     * 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 ProductionExceptionResponseResponse {

        private ProductionExceptionHandlerResponse productionExceptionHandlerResponse;

        private List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords;

        /**
         * Constructs a new {@code ProductionExceptionResponseResponse} 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 ProductionExceptionResponseResponse(final ProductionExceptionHandlerResponse productionExceptionHandlerResponse,
                                            final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            this.productionExceptionHandlerResponse = productionExceptionHandlerResponse;
            this.deadLetterQueueRecords = deadLetterQueueRecords;
        }

        /**
         * Creates a {@code ProductionExceptionResponseResponse} 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 ProductionExceptionResponseResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#FAIL} status.
         */
        public static ProductionExceptionResponseResponse failProcessingfail(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new ProductionExceptionResponseResponse(ProductionExceptionHandlerResponse.FAIL, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code ProductionExceptionResponseResponse} indicating that processing should fail.
         *
         * @return a {@code ProductionExceptionResponseResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#FAIL} status.
         */
        public static ProductionExceptionResponseResponse failProcessingfail() {
            return failProcessingfail(Collections.emptyList());
        }

        /**
         * Creates a {@code ProductionExceptionResponseResponse} 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 ProductionExceptionResponseResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
         */
        public static ProductionExceptionResponseResponse continueProcessingresume(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new ProductionExceptionResponseResponse(ProductionExceptionHandlerResponse.CONTINUE, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code ProductionExceptionResponseResponse} indicating that processing should continue.
         *
         * @return a {@code ProductionExceptionResponseResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
         */
        public static ProductionExceptionResponseResponse continueProcessingresume() {
            return continueProcessingresume(Collections.emptyList());
            }

        /**
         * Creates a {@code ProductionExceptionResponseResponse} indicating that processing should retry.
         *
         * @return a {@code ProductionExceptionResponseResponse} with a {@link DeserializationExceptionHandler.DeserializationHandlerResponse#CONTINUE} status.
         */
        public static ProductionExceptionResponseResponse retryProcessingretry() {
            return new ProductionExceptionResponseResponse(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);
        }
    }
      ...  
}

DeserializationExceptionHandler.java

Changes:

  • Adding a DeserializationExceptionResponse Response nested class,that contains a DeserializationHandlerResponse indicating whether to continue processing or fail, and the list of records to be sent to the dead letter queue topic
  • Deprecating the handler() method and adding a new method: handlerError(). The return type is the DeserializationExceptionResponse
Code Block
languagejava
public interface DeserializationExceptionHandler extends Configurable {

   ...      
  
    @Deprecated
    default DeserializationHandlerResponse handle(final ErrorHandlerContext context,
                                                  final ConsumerRecord<byte[], byte[]> record,
                                                  final Exception exception) {
        throw new UnsupportedOperationException();
    }

    /**
     * Inspects a record and the exception received during deserialization.
     *
     * @param context
     *     Error handler context.
     * @param record
     *     Record that failed deserialization.
     * @param exception
     *     The actual exception.
     *
     * @return a {@link DeserializationExceptionResponse} object
     */
    default DeserializationExceptionResponse handleError(final ErrorHandlerContext context, final ConsumerRecord<byte[], byte[]> record, final Exception exception) {           return new DeserializationExceptionResponse(handle(context, record, exception),   Collections.emptyList());
    }     

    ...

     /**
     * Represents the result of handling a deserialization exception.
     * <p>
     * The {@code Response} class encapsulates a {@link ProcessingExceptionHandler.ProcessingHandlerResponse},
     * 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 DeserializationExceptionResponseResponse {

        private DeserializationHandlerResponse deserializationHandlerResponse;

        private List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords;

        /**
         * Constructs a new {@code DeserializationExceptionResponse} object.
         *
         * @param deserializationHandlerResponse 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 DeserializationExceptionResponseResponse(final DeserializationHandlerResponse deserializationHandlerResponse,
                                                 final List<ProducerRecord<byte[final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            this.deserializationHandlerResponse = deserializationHandlerResponse;
            this.deadLetterQueueRecords = deadLetterQueueRecords;
        }

        /**
         * Creates a {@code DeserializationExceptionResponseResponse} 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 DeserializationExceptionResponseResponse} with a {@link DeserializationHandlerResponse#FAIL} status.
         */
        public static DeserializationExceptionResponseResponse failProcessingfail(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new DeserializationExceptionResponseResponse(DeserializationHandlerResponse.FAIL, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code DeserializationExceptionResponseResponse} indicating that processing should fail.
         *
         * @return a {@code DeserializationExceptionResponseResponse} with a {@link DeserializationHandlerResponse#FAIL} status.
         */
        public static DeserializationExceptionResponseResponse failProcessingfail() {
            return failProcessingfail(Collections.emptyList());
        }

        /**
         * Creates a {@code DeserializationExceptionResponseResponse} 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 DeserializationExceptionResponseResponse} with a {@link DeserializationHandlerResponse#CONTINUE} status.
         */
        public static DeserializationExceptionResponseResponse continueProcessingresume(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new DeserializationExceptionResponseResponse(DeserializationHandlerResponse.CONTINUE, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code DeserializationExceptionResponseResponse} indicating that processing should continue.
         *
         * @return a {@code DeserializationExceptionResponseResponse} with a {@link DeserializationHandlerResponse#CONTINUE} status.
         */
        public static DeserializationExceptionResponseResponse continueProcessingresume() {
            return continueProcessingresume(Collections.emptyList());
        }

        /**
         * Retrieves the deserialization handler response.
         *
         * @return the {@link DeserializationHandlerResponse} indicating whether processing should continue or fail.
         */
        public DeserializationHandlerResponse response() {
            return deserializationHandlerResponse;
        }

        /**
         * 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);
        }
    }

   

}

ProcessingExceptionHandler.java

Changes:

  • Adding a ProcessingExceptionResponse Response nested class,that contains a ProcessingHandlerResponse indicating whether to continue processing or fail, and the list of records to be sent to the dead letter queue topic
  • Deprecating the handler() method and adding a new method: handlerError(). The return type is the ProcessingExceptionResponse
Code Block
languagejava
public interface ProcessingExceptionHandler extends Configurable {

   ...   
    
    @Deprecated
    default ProcessingHandlerResponse handle(final ErrorHandlerContext context, final Record<?, ?> record, final Exception exception){
        throw new UnsupportedOperationException();
    };

    /**
     * Inspects a record and the exception received during processing.
     *
     * @param context
     *     Processing context metadata.
     * @param record
     *     Record where the exception occurred.
     * @param exception
     *     The actual exception.
     *
     * @return a {@link ProcessingExceptionResponse} object
     */
    default ProcessingExceptionResponse handleError(final ErrorHandlerContext context, final Record<?, ?> record, final Exception exception) {                 
       return new ProcessingExceptionResponse(handle(context, record, exception), Collections.emptyList());
    }   

...

     /**
     * Represents the result of handling a processing exception.
     * <p>
     * The {@code Response} class encapsulates a {@link ProcessingHandlerResponse},
     * indicating whether processing should continue or fail, along with an optional list of
     * {@link org.apache.kafka.clients.producer.ProducerRecord} instances to be sent to a dead letter queue.
     * </p>
     */
    class ProcessingExceptionResponseResponse {

        private ProcessingHandlerResponse processingHandlerResponse;

        private List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords;

        /**
         * Constructs a new {@code ProcessingExceptionResponse} object.
         *
         * @param processingHandlerResponse 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 ProcessingExceptionResponseResponse(final ProcessingHandlerResponse processingHandlerResponse,
                                            final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            this.processingHandlerResponse = processingHandlerResponse;
            this.deadLetterQueueRecords = deadLetterQueueRecords;
        }

        /**
         * Creates a {@code ProcessingExceptionResponseResponse} 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 ProcessingExceptionResponse} with a {@link ProcessingHandlerResponse#FAIL} status.
         */
        public static ProcessingExceptionResponseResponse failProcessingfail(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new ProcessingExceptionResponseResponse(ProcessingHandlerResponse.FAIL, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code ProcessingExceptionResponseResponse} indicating that processing should fail.
         *
         * @return a {@code ProcessingExceptionResponseResponse} with a {@link ProcessingHandlerResponse#FAIL} status.
         */
        public static ProcessingExceptionResponseResponse failProcessingfail() {
            return failProcessingfail(Collections.emptyList());
        }

        /**
         * Creates a {@code ProcessingExceptionResponseResponse} 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 Response} with a {@link ProcessingHandlerResponse#CONTINUE} status.
         */
        public static ProcessingExceptionResponseResponse continueProcessingresume(final List<ProducerRecord<byte[], byte[]>> deadLetterQueueRecords) {
            return new ProcessingExceptionResponseResponse(ProcessingHandlerResponse.CONTINUE, deadLetterQueueRecords);
        }

        /**
         * Creates a {@code ProcessingExceptionResponseResponse} indicating that processing should continue.
         *
         * @return a {@code ProcessingExceptionResponseResponse} with a {@link ProcessingHandlerResponse#CONTINUE} status.
         */
        public static ProcessingExceptionResponseResponse continueProcessingresume() {
            return continueProcessingresume(Collections.emptyList());
        }

        /**
         * Retrieves the processing handler response.
         *
         * @return the {@link ProcessingHandlerResponse} indicating whether processing should continue or fail.
         */
        public ProcessingHandlerResponse response() {
            return processingHandlerResponse;
        }

        /**
         * 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);
        }
    }
  

}

RecordContext.java

Changes:

...