Versions Compared

Key

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

Table of Contents

This page is meant as a template for writing a KIP. To create a KIP choose Tools->Copy on this page and modify with your content and replace the heading with the next KIP number and a description of your issue. Replace anything in italics with your own description.

Status

Current state: Under Discussion Accepted

Discussion thread: here [Change the link from the KIP proposal email archive to your own email thread] 

Vote thread: here 

JIRA:

Jira
serverASF JIRA
serverId5aa69414-a9e9-3523-82ec-879b028fb15b
keyKAFKA-10790

...

Public Interfaces

This KIP proposes two ways a way to implement this enhancement:

...

If the `Producer` detects that the `flush` method is invoked within the callback of the `close` method, throw a `KafkaException`.

This approach would require modifying org.apache.kafka.clients.producer.Producer#flush  to include KafkaException .

Proposed Changes

If we choose to follow point 2, we will need to make the following changesThe following code is a demonstration of how we will add to the flush method:

Code Block
languagejava
titleProducer#flush
    /**
	 *                Omitting the unchange section ......
     * <p>
     * <b>Important:</b> This method must not be called from within the callback provided to
     * {@link #send(ProducerRecord, Callback)}.Invoking <code>flush()</code> in this context will result in a
     * {@link KafkaException} being thrown, as it will cause a deadlock.
     * </p>
     *
     * @throws InterruptException If the thread is interrupted while blocked
     * @throws KafkaException If the method is invoked inside a {@link #send(ProducerRecord, Callback)} callback
     */
    @Override
    public void flush() {
        if (Thread.currentThread() == this.ioThread) {
            log.error("KafkaProducer.flush() invocation inside a callback is not permitted because it may lead to deadlock.");
            throw new KafkaException("KafkaProducer.flush() invocation inside a callback is not permitted because it may lead to deadlock.");
        }
		
		// Omitting the unchange section ......

    }

...

Add a unit test to ensure that the deadlock protection throws a `KafkaException` with the corresponding error message to inform the user.

Rejected Alternatives

Adding an error log:

If the `Producer` detects that the `flush` method is invoked within the callback of the `close` method, log an error message.

Reason for Rejection

Although adding an error log message could help users understand what happened, it does not provide sufficient protection.

This KIP aims to offer fast-fail protection, which would save users time by allowing the `send` method to return immediately.N/A