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

  1. 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.
  2. If a mock SinkTaskContext is in play, configure it to expose a consumer or avoid code paths that need consumerGroupMetadata.
  3. If using an alternative Connect-compatible runtime, patch KafkaUtils to support its context implementation.
  4. 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

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


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)