{"record":{"id":"139199f163bca14a","repo":"apache/kafka","slug":"failed-to-construct-kafka-consumer","errorCode":null,"errorMessage":"Failed to construct kafka consumer","messagePattern":"Failed to construct kafka consumer","errorType":"exception","errorClass":"KafkaException","httpStatus":null,"severity":"critical","filePath":"clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java","lineNumber":602,"sourceCode":"                    deserializers,\n                    fetchMetricsManager,\n                    time);\n\n            if (groupMetadata.get().isPresent() &&\n                GroupProtocol.of(config.getString(ConsumerConfig.GROUP_PROTOCOL_CONFIG)) == GroupProtocol.CONSUMER) {\n                config.ignore(ConsumerConfig.GROUP_REMOTE_ASSIGNOR_CONFIG); // Used by background thread\n            }\n            config.logUnused();\n            AppInfoParser.registerAppInfo(CONSUMER_JMX_PREFIX, clientId, metrics, time.milliseconds());\n            log.debug(\"Kafka consumer initialized\");\n        } catch (Throwable t) {\n            // call close methods if internal objects are already constructed; this is to prevent resource leak. see KAFKA-2121\n            // we do not need to call `close` at all when `log` is null, which means no internal objects were initialized.\n            if (this.log != null) {\n                close(Duration.ZERO, CloseOptions.GroupMembershipOperation.LEAVE_GROUP, true);\n            }\n            // now propagate the exception\n            throw new KafkaException(\"Failed to construct kafka consumer\", t);\n        }\n    }\n\n    // Visible for testing\n    AsyncKafkaConsumer(LogContext logContext,\n                       String clientId,\n                       Deserializers<K, V> deserializers,\n                       FetchBuffer fetchBuffer,\n                       FetchCollector<K, V> fetchCollector,\n                       FetchMetricsManager fetchMetricsManager,\n                       RebalanceCallbackMetricsManager rebalanceCallbackMetricsManager,\n                       ConsumerInterceptors<K, V> interceptors,\n                       Time time,\n                       ApplicationEventHandler applicationEventHandler,\n                       BlockingQueue<BackgroundEvent> backgroundEventQueue,\n                       CompletableEventReaper backgroundEventReaper,\n                       ConsumerRebalanceListenerInvoker rebalanceListenerInvoker,\n                       Metrics metrics,","sourceCodeStart":584,"sourceCodeEnd":620,"githubUrl":"https://github.com/apache/kafka/blob/996fb4585aa1bcc8980b0e1b8d6b168b986cd979/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java#L584-L620","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// before\nProperties p = new Properties();\np.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, \"localhost:9092\");\np.put(ConsumerConfig.GROUP_ID_CONFIG, \"g\");\nnew KafkaConsumer<>(p, null, null); // -> KafkaException: Failed to construct kafka consumer (cause: missing deserializer)\n\n// after\nProperties p = new Properties();\np.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, \"localhost:9092\");\np.put(ConsumerConfig.GROUP_ID_CONFIG, \"g\");\nnew KafkaConsumer<>(p, new StringDeserializer(), new StringDeserializer());","handlingStrategy":"try-catch","validationCode":"// Pre-validate required configs before constructing.\nObjects.requireNonNull(props.getProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG), \"key.deserializer required\");\nObjects.requireNonNull(props.getProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG), \"value.deserializer required\");","typeGuard":"null","tryCatchPattern":"try {\n    this.consumer = new KafkaConsumer<>(props, keyDeser, valueDeser);\n} catch (KafkaException e) {\n    Throwable cause = e.getCause();\n    log.error(\"consumer construction failed: {}\", cause == null ? e : cause);\n    throw e;\n}","preventionTips":["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."],"tags":["consumer","construction","config","kafka-clients"],"backgroundTag":null,"analyzedSha":"996fb4585aa1bcc8980b0e1b8d6b168b986cd979","analyzedAt":"2026-08-11T22:03:28.655Z","contentChangedAt":null,"schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}