apache/iceberg · error · ConnectException

Cannot retrieve members for consumer group

Error message

Cannot retrieve members for consumer group: ${consumerGroupId}

What it means

KafkaUtils.consumerGroupDescription wraps failures from the AdminClient describeConsumerGroups call into a ConnectException. The sink needs the consumer group's members to map topic partitions to committer tasks; if the describe call is interrupted or fails asynchronously, this error reports it.

Solutions

  1. Verify the consumer group exists: kafka-consumer-groups --bootstrap-server ... --describe --group <groupId>.
  2. Check worker ACLs — the principal needs DESCRIBE on the consumer group resource.
  3. Verify bootstrap servers and network connectivity from the Connect worker to the brokers.
  4. If the group was recently recreated, restart the connector after group metadata is stable.
Defensive patterns

Strategy: try-catch

Validate before calling

// before starting the connector, verify group access
kafka-consumer-groups --bootstrap-server broker:9092 --describe --group <consumer.group.id>

Try / catch

try {
  ConsumerGroupDescription desc = KafkaUtils.consumerGroupDescription(admin, groupId);
} catch (ConnectException e) {
  LOG.error("Failed to describe consumer group {} — check existence, ACLs (DESCRIBE), broker connectivity", groupId, e);
  throw e;
}

Prevention

When it happens

Trigger: admin.describeConsumerGroups(consumerGroupId).describedGroups().get(...).get() throws InterruptedException or ExecutionException — group does not exist, insufficient ACLs, broker unreachable, or the future was interrupted.

Common situations: Wrong consumer.group.id or group deleted/recreated between task start; Connect worker lacking DESCRIBE ACL on the group; network/firewall issue to the Kafka brokers; group metadata not yet propagated at startup.

Understand the failure class

Background: Record Not Found Errors: "not found", RecordNotFound, and "was not found" — what they mean and how to fix them — this error's family across 28 libraries.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/1512d095188f62d9. Report an issue: GitHub.

Appendix: source

Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java:51

import org.apache.kafka.connect.sink.SinkTaskContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

class KafkaUtils {

  private static final Logger LOG = LoggerFactory.getLogger(KafkaUtils.class);

  private static final String CONTEXT_CLASS_NAME =
      "org.apache.kafka.connect.runtime.WorkerSinkTaskContext";

  static ConsumerGroupDescription consumerGroupDescription(String consumerGroupId, Admin admin) {
    try {
      DescribeConsumerGroupsResult result =
          admin.describeConsumerGroups(ImmutableList.of(consumerGroupId));
      return result.describedGroups().get(consumerGroupId).get();

    } catch (InterruptedException | ExecutionException e) {
      throw new ConnectException(
          "Cannot retrieve members for consumer group: " + consumerGroupId, e);
    }
  }

  static ConsumerGroupMetadata consumerGroupMetadata(SinkTaskContext context) {
    return kafkaConsumer(context).groupMetadata();
  }

  static void seekToLastCommittedOffsets(SinkTaskContext context) {
    Consumer<byte[], byte[]> consumer = kafkaConsumer(context);
    if (consumer == null) {
      return;
    }

    Map<TopicPartition, OffsetAndMetadata> committedOffsets =
        consumer.committed(consumer.assignment());
    if (committedOffsets == null || committedOffsets.isEmpty()) {
      return;

View on GitHub (pinned to 86d9c8fc54)