apache/iceberg · error · ConnectException
Unable to retrieve consumer from context
Error message
Unable to retrieve consumer from context: ${contextClassName} What it means
KafkaUtils.kafkaConsumer uses reflection (DynFields on the hidden 'consumer' field of the SinkTaskContext implementation class) to reach the underlying Kafka Consumer because no public API exposes it. Any failure — class renamed, field moved, reflective access blocked — becomes this ConnectException naming the context class.
Solutions
- Check that the connector runs on a real Kafka Connect worker with the Kafka version the Iceberg sink was built for — align kafka-clients versions.
- If a mock SinkTaskContext is in play, configure it to expose a consumer or avoid code paths that need consumerGroupMetadata.
- If using an alternative Connect-compatible runtime, patch KafkaUtils to support its context implementation.
- Open/add reflective-access allowances (--add-opens) if module encapsulation is the blocker.
Defensive patterns
Strategy: try-catch
Validate before calling
// ensure a real Kafka context is in use before reflective access
boolean realContext = context.getClass().getName().startsWith("org.apache.kafka.connect.runtime."); Try / catch
try {
Consumer<byte[], byte[]> consumer = KafkaUtils.consumer(context);
} catch (ConnectException e) {
LOG.error("Reflective consumer lookup failed for context {} — check Kafka version alignment / mock context", e.getMessage(), e);
throw e;
} Prevention
- Pin kafka-clients versions to those the Iceberg sink was built against
- Do not run production paths against mock SinkTaskContext implementations
- Add --add-opens reflective access if running under JPMS strong encapsulation
- Pin kafka-clients versions to those the Iceberg sink was built against
When it happens
Trigger: The runtime SinkTaskContext implementation is not the expected Kafka class (CONTEXT_CLASS_NAME mismatch, e.g. mock/testing framework or alternative Connect implementation), or a Kafka upgrade removed/renamed the hidden 'consumer' field, or a SecurityManager/module access rule blocks reflective access.
Common situations: Running sink unit/integration tests with a mock SinkTaskContext; Kafka client major-version upgrade changing WorkerSinkTask internals; JPMS strong encapsulation blocking DynFields hidden access.
Understand the failure class
Background: "is not a compatible type" / "cannot merge" errors: when a value's type doesn't match what the library requires — this error's family across 65 libraries.
Related errors
- An error occurred closing catalog instance, ignoring...
- An error occurred converting record, topic
- Cannot bind constructors
- Cannot convert date
- Cannot convert java.util.Date to variant without a…
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/426bde8394a3b459.
Report an issue: GitHub.
Appendix: source
Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java:94
consumer.seek(topicPartition, offsetAndMetadata.offset());
} catch (IllegalStateException e) {
LOG.warn(
"Rebalance may have occurred, partition {} lost before seeking",
topicPartition,
e);
}
}
});
}
@SuppressWarnings("unchecked")
private static Consumer<byte[], byte[]> kafkaConsumer(SinkTaskContext context) {
String contextClassName = context.getClass().getName();
try {
return ((Consumer<byte[], byte[]>)
DynFields.builder().hiddenImpl(CONTEXT_CLASS_NAME, "consumer").build(context).get());
} catch (Exception e) {
throw new ConnectException(
"Unable to retrieve consumer from context: " + contextClassName, e);
}
}
private KafkaUtils() {}
}
View on GitHub (pinned to 86d9c8fc54)