diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterConsumerGroupOffsetsHandler.java b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterConsumerGroupOffsetsHandler.java index eab2e2bb73a40..425ed66bd29a2 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterConsumerGroupOffsetsHandler.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterConsumerGroupOffsetsHandler.java @@ -179,6 +179,8 @@ private void handleError( case INVALID_GROUP_ID: case INVALID_COMMIT_OFFSET_SIZE: case GROUP_AUTHORIZATION_FAILED: + // Member level errors. + case UNKNOWN_MEMBER_ID: log.debug("OffsetCommit request for group id {} failed due to error {}.", groupId.idValue, error); partitionResults.put(topicPartition, error); diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointTask.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointTask.java index 959961812ea5a..3e6247334bb81 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointTask.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointTask.java @@ -17,9 +17,11 @@ package org.apache.kafka.connect.mirror; import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.AlterConsumerGroupOffsetsResult; import org.apache.kafka.clients.admin.ConsumerGroupDescription; import org.apache.kafka.common.ConsumerGroupState; import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.errors.UnknownMemberIdException; import org.apache.kafka.connect.source.SourceTask; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.data.Schema; @@ -306,9 +308,18 @@ Map> syncGroupOffset() { void syncGroupOffset(String consumerGroupId, Map offsetToSync) { if (targetAdminClient != null) { - targetAdminClient.alterConsumerGroupOffsets(consumerGroupId, offsetToSync); - log.trace("sync-ed the offset for consumer group: {} with {} number of offset entries", - consumerGroupId, offsetToSync.size()); + AlterConsumerGroupOffsetsResult result = targetAdminClient.alterConsumerGroupOffsets(consumerGroupId, offsetToSync); + result.all().whenComplete((v, throwable) -> { + if (throwable != null) { + if (throwable.getCause() instanceof UnknownMemberIdException) { + log.warn("Unable to sync offsets for consumer group {}. This is likely caused by consumers currently using this group in the target cluster.", consumerGroupId); + } else { + log.error("Unable to sync offsets for consumer group {}.", consumerGroupId, throwable); + } + } else { + log.trace("Sync-ed {} offsets for consumer group {}.", offsetToSync.size(), consumerGroupId); + } + }); } }