From e0f14dfd6552eac18e7c3c000e3622c1d5f37fa7 Mon Sep 17 00:00:00 2001 From: Onur Karaman Date: Fri, 18 Sep 2015 15:35:10 -0700 Subject: [PATCH] separate REBALANCE_IN_PROGRESS and ILLEGAL_GENERATION error codes --- .../kafka/clients/consumer/internals/Coordinator.java | 4 ++++ .../main/java/org/apache/kafka/common/protocol/Errors.java | 4 +++- .../main/scala/kafka/coordinator/ConsumerCoordinator.scala | 4 +++- .../kafka/coordinator/ConsumerCoordinatorResponseTest.scala | 6 +++--- 4 files changed, 13 insertions(+), 5 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Coordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Coordinator.java index e7ffe252201fa..ee389692d4e6c 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Coordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/Coordinator.java @@ -641,6 +641,10 @@ public void handle(HeartbeatResponse heartbeatResponse, RequestFuture futu log.info("Attempt to heart beat failed since coordinator is either not started or not valid, marking it as dead."); coordinatorDead(); future.raise(Errors.forCode(error)); + } else if (error == Errors.REBALANCE_IN_PROGRESS.code()) { + log.info("Attempt to heart beat failed since the group is rebalancing, try to re-join group."); + subscriptions.needReassignment(); + future.raise(Errors.REBALANCE_IN_PROGRESS); } else if (error == Errors.ILLEGAL_GENERATION.code()) { log.info("Attempt to heart beat failed since generation id is not legal, try to re-join group."); subscriptions.needReassignment(); diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java b/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java index b3415c3a8113f..220132f7f023d 100644 --- a/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java +++ b/clients/src/main/java/org/apache/kafka/common/protocol/Errors.java @@ -88,7 +88,9 @@ public enum Errors { new ApiException("Some of the committing partitions are not assigned the committer")), INVALID_COMMIT_OFFSET_SIZE(28, new ApiException("The committing offset data size is not valid")), - AUTHORIZATION_FAILED(29, new ApiException("Request is not authorized.")); + AUTHORIZATION_FAILED(29, new ApiException("Request is not authorized.")), + REBALANCE_IN_PROGRESS(30, + new ApiException("The group is rebalancing, so a rejoin is needed.")); private static final Logger log = LoggerFactory.getLogger(Errors.class); diff --git a/core/src/main/scala/kafka/coordinator/ConsumerCoordinator.scala b/core/src/main/scala/kafka/coordinator/ConsumerCoordinator.scala index 1bceb43793a38..64e21c5b0fc8b 100644 --- a/core/src/main/scala/kafka/coordinator/ConsumerCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/ConsumerCoordinator.scala @@ -210,8 +210,10 @@ class ConsumerCoordinator(val brokerId: Int, responseCallback(Errors.UNKNOWN_CONSUMER_ID.code) } else if (!group.has(consumerId)) { responseCallback(Errors.UNKNOWN_CONSUMER_ID.code) - } else if (generationId != group.generationId || !group.is(Stable)) { + } else if (generationId != group.generationId) { responseCallback(Errors.ILLEGAL_GENERATION.code) + } else if (!group.is(Stable)) { + responseCallback(Errors.REBALANCE_IN_PROGRESS.code) } else { val consumer = group.get(consumerId) completeAndScheduleNextHeartbeatExpiration(group, consumer) diff --git a/core/src/test/scala/unit/kafka/coordinator/ConsumerCoordinatorResponseTest.scala b/core/src/test/scala/unit/kafka/coordinator/ConsumerCoordinatorResponseTest.scala index 42ffddec2104f..07f7326d72bf2 100644 --- a/core/src/test/scala/unit/kafka/coordinator/ConsumerCoordinatorResponseTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/ConsumerCoordinatorResponseTest.scala @@ -232,7 +232,7 @@ class ConsumerCoordinatorResponseTest extends JUnitSuite { } @Test - def testHeartbeatDuringRebalanceCausesIllegalGeneration() { + def testHeartbeatDuringRebalanceCausesRebalanceInProgress() { val groupId = "groupId" val partitionAssignmentStrategy = "range" @@ -249,10 +249,10 @@ class ConsumerCoordinatorResponseTest extends JUnitSuite { sendJoinGroup(groupId, JoinGroupRequest.UNKNOWN_CONSUMER_ID, partitionAssignmentStrategy, DefaultSessionTimeout, isCoordinatorForGroup = true) - // We should be in the middle of a rebalance, so the heartbeat should return illegal generation + // We should be in the middle of a rebalance, so the heartbeat should return rebalance in progress EasyMock.reset(offsetManager) val heartbeatResult = heartbeat(groupId, assignedConsumerId, initialGenerationId, isCoordinatorForGroup = true) - assertEquals(Errors.ILLEGAL_GENERATION.code, heartbeatResult) + assertEquals(Errors.REBALANCE_IN_PROGRESS.code, heartbeatResult) } @Test