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
- Inspect KafkaException.getCause() (and its cause chain) - the real error (e.g. ConfigException, ClassNotFoundException) is what to fix.
- Ensure key.deserializer and value.deserializer are set and point to installed classes (or use the String/Long/ByteBuffer overloads of KafkaConsumer).
- Validate security config (ssl.truststore.location, sasl.jaas.config) for missing files / unset credentials.
- 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
- Always pass deserializer instances to the constructor rather than relying on class-name config.
- Log KafkaException.getCause() so the real failure is visible.
- Validate security config (truststore path, jaas) before construction.
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
- The configured group.id should not be an empty string or…
- cannot be set when using a share group.
- Consumer is not subscribed to any topics or assigned any…
- enable.auto.commit cannot be set to true when default group…
- Invalid negative offset
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)