apache/kafka · critical · IOException

Client was shutdown before response was read

Error message

Client was shutdown before response was read

What it means

Thrown by NetworkClientUtils.sendAndReceive() when the poll loop exits because client.active() returns false — the NetworkClient was shut down (state left ACTIVE) before the response for the request could be read. This is a terminal failure: the client is no longer usable.

Solutions

  1. Ensure the client is not closed while sendAndReceive is in progress — coordinate lifecycle with a flag or latch.
  2. Register a shutdown hook that waits for in-flight operations to complete before closing.
  3. Avoid sharing a single NetworkClient across components with independent close behavior.
Defensive patterns

Strategy: try-catch

Validate before calling

if (!client.active()) {
    throw new IllegalStateException("Client inactive; cannot sendAndReceive");
}

Try / catch

try {
    NetworkClientUtils.sendAndReceive(client, request, time);
} catch (IOException e) {
    if (!client.active()) {
        // shutdown race: do not retry, propagate
    }
}

Prevention

When it happens

Trigger: After sending a request, the while(client.active()) loop terminates because the client was closed (or is closing). The response never arrived before shutdown completed.

Common situations: Calling close() on the client from another thread while sendAndReceive is blocking, a JVM shutdown hook closing clients during in-flight requests, or application logic that closes the client before draining pending operations.

Related errors


AI-assisted analysis of apache/kafka@996fb4585a (2026-08-11). Data as JSON: /api/errors/0ddc90b085b2b1bd. Report an issue: GitHub.

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/NetworkClientUtils.java:120

     */
    public static ClientResponse sendAndReceive(KafkaClient client, ClientRequest request, Time time) throws IOException {
        try {
            client.send(request, time.milliseconds());
            while (client.active()) {
                List<ClientResponse> responses = client.poll(Long.MAX_VALUE, time.milliseconds());
                for (ClientResponse response : responses) {
                    if (response.requestHeader().correlationId() == request.correlationId()) {
                        if (response.wasDisconnected()) {
                            throw new IOException("Connection to " + response.destination() + " was disconnected before the response was read");
                        }
                        if (response.versionMismatch() != null) {
                            throw response.versionMismatch();
                        }
                        return response;
                    }
                }
            }
            throw new IOException("Client was shutdown before response was read");
        } catch (DisconnectException e) {
            if (client.active())
                throw e;
            else
                throw new IOException("Client was shutdown before response was read");

        }
    }

    /**
     * Check if the code is disconnected and unavailable for immediate reconnection (i.e. if it is in
     * reconnect backoff window following the disconnect).
     */
    public static boolean isUnavailable(KafkaClient client, Node node, Time time) {
        return client.connectionFailed(node) && client.connectionDelay(node, time.milliseconds()) > 0;
    }

    /**

View on GitHub (pinned to 996fb4585a)