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
- Verify the consumer group exists: kafka-consumer-groups --bootstrap-server ... --describe --group <groupId>.
- Check worker ACLs — the principal needs DESCRIBE on the consumer group resource.
- Verify bootstrap servers and network connectivity from the Connect worker to the brokers.
- 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
- Verify the consumer group exists and is stable before connector start
- Grant DESCRIBE ACL on the group to the Connect worker principal
- Validate broker connectivity from the worker with a smoke test
- Avoid deleting/recreating the group while the connector is running
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
- An error occurred closing catalog instance, ignoring...
- An error occurred converting record, topic
- Cannot convert date
- Cannot convert java.util.Date to variant without a…
- Cannot convert map to variant: keys must be non-null…
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)