{"record":{"id":"426bde8394a3b459","repo":"apache/iceberg","slug":"unable-to-retrieve-consumer-from-context-contex","errorCode":null,"errorMessage":"Unable to retrieve consumer from context: ${contextClassName}","messagePattern":"Unable to retrieve consumer from context: (.+?)","errorType":"exception","errorClass":"ConnectException","httpStatus":null,"severity":"error","filePath":"kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java","lineNumber":94,"sourceCode":"              consumer.seek(topicPartition, offsetAndMetadata.offset());\n            } catch (IllegalStateException e) {\n              LOG.warn(\n                  \"Rebalance may have occurred, partition {} lost before seeking\",\n                  topicPartition,\n                  e);\n            }\n          }\n        });\n  }\n\n  @SuppressWarnings(\"unchecked\")\n  private static Consumer<byte[], byte[]> kafkaConsumer(SinkTaskContext context) {\n    String contextClassName = context.getClass().getName();\n    try {\n      return ((Consumer<byte[], byte[]>)\n          DynFields.builder().hiddenImpl(CONTEXT_CLASS_NAME, \"consumer\").build(context).get());\n    } catch (Exception e) {\n      throw new ConnectException(\n          \"Unable to retrieve consumer from context: \" + contextClassName, e);\n    }\n  }\n\n  private KafkaUtils() {}\n}\n","sourceCodeStart":76,"sourceCodeEnd":101,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java#L76-L101","documentation":"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.","triggerScenarios":"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.","commonSituations":"Running sink unit/integration tests with a mock SinkTaskContext; Kafka client major-version upgrade changing WorkerSinkTask internals; JPMS strong encapsulation blocking DynFields hidden access.","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."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// ensure a real Kafka context is in use before reflective access\nboolean realContext = context.getClass().getName().startsWith(\"org.apache.kafka.connect.runtime.\");","typeGuard":null,"tryCatchPattern":"try {\n  Consumer<byte[], byte[]> consumer = KafkaUtils.consumer(context);\n} catch (ConnectException e) {\n  LOG.error(\"Reflective consumer lookup failed for context {} — check Kafka version alignment / mock context\", e.getMessage(), e);\n  throw e;\n}","preventionTips":["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"],"tags":["kafka-connect","reflection","version-incompatibility","sinktaskcontext"],"backgroundTag":"incompatible-source-type","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}