apache/kafka · critical · KafkaException

Failed to construct kafka consumer

Error message

Failed to construct kafka consumer

What it means

AsyncKafkaConsumer's constructor wraps any Throwable thrown during initialization in KafkaException('Failed to construct kafka consumer', t) after attempting to close already-built internals (KAFKA-2121 resource-leak guard). The original failure is the cause (KafkaException.getCause()); the message itself is generic. Common root causes include missing/wrong deserializers, invalid bootstrap servers, missing security credentials, and unsupported config combinations.

Solutions

  1. Inspect KafkaException.getCause() (and its cause chain) - the real error (e.g. ConfigException, ClassNotFoundException) is what to fix.
  2. Ensure key.deserializer and value.deserializer are set and point to installed classes (or use the String/Long/ByteBuffer overloads of KafkaConsumer).
  3. Validate security config (ssl.truststore.location, sasl.jaas.config) for missing files / unset credentials.
  4. Cross-check every config key against ConsumerConfig names to catch typos that surface as ConfigException.

Example fix

// before
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
p.put(ConsumerConfig.GROUP_ID_CONFIG, "g");
new KafkaConsumer<>(p, null, null); // -> KafkaException: Failed to construct kafka consumer (cause: missing deserializer)

// after
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
p.put(ConsumerConfig.GROUP_ID_CONFIG, "g");
new KafkaConsumer<>(p, new StringDeserializer(), new StringDeserializer());
Defensive patterns

Strategy: try-catch

Validate before calling

// Pre-validate required configs before constructing.
Objects.requireNonNull(props.getProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG), "key.deserializer required");
Objects.requireNonNull(props.getProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG), "value.deserializer required");

Type guard

null

Try / catch

try {
    this.consumer = new KafkaConsumer<>(props, keyDeser, valueDeser);
} catch (KafkaException e) {
    Throwable cause = e.getCause();
    log.error("consumer construction failed: {}", cause == null ? e : cause);
    throw e;
}

Prevention

When it happens

Trigger: new KafkaConsumer<>(props) where props lacks key.deserializer/value.deserializer, has an invalid value for a typed config (ConfigException surfaced as the cause), references a security scheme whose classes are missing, or fails serializer instantiation.

Common situations: Running with a properties file that omits deserializers; passing Strings where LongDeserializer is needed; SSL/SASL config referencing keystore paths that do not exist; classpath missing a custom serializer dependency.

Related errors


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

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java:602

                    deserializers,
                    fetchMetricsManager,
                    time);

            if (groupMetadata.get().isPresent() &&
                GroupProtocol.of(config.getString(ConsumerConfig.GROUP_PROTOCOL_CONFIG)) == GroupProtocol.CONSUMER) {
                config.ignore(ConsumerConfig.GROUP_REMOTE_ASSIGNOR_CONFIG); // Used by background thread
            }
            config.logUnused();
            AppInfoParser.registerAppInfo(CONSUMER_JMX_PREFIX, clientId, metrics, time.milliseconds());
            log.debug("Kafka consumer initialized");
        } catch (Throwable t) {
            // call close methods if internal objects are already constructed; this is to prevent resource leak. see KAFKA-2121
            // we do not need to call `close` at all when `log` is null, which means no internal objects were initialized.
            if (this.log != null) {
                close(Duration.ZERO, CloseOptions.GroupMembershipOperation.LEAVE_GROUP, true);
            }
            // now propagate the exception
            throw new KafkaException("Failed to construct kafka consumer", t);
        }
    }

    // Visible for testing
    AsyncKafkaConsumer(LogContext logContext,
                       String clientId,
                       Deserializers<K, V> deserializers,
                       FetchBuffer fetchBuffer,
                       FetchCollector<K, V> fetchCollector,
                       FetchMetricsManager fetchMetricsManager,
                       RebalanceCallbackMetricsManager rebalanceCallbackMetricsManager,
                       ConsumerInterceptors<K, V> interceptors,
                       Time time,
                       ApplicationEventHandler applicationEventHandler,
                       BlockingQueue<BackgroundEvent> backgroundEventQueue,
                       CompletableEventReaper backgroundEventReaper,
                       ConsumerRebalanceListenerInvoker rebalanceListenerInvoker,
                       Metrics metrics,

View on GitHub (pinned to 996fb4585a)