From cfd8e8e26ad75bafb94094e2b76acf4a330c4508 Mon Sep 17 00:00:00 2001 From: Jiangjie Qin Date: Fri, 18 Sep 2015 10:21:27 -0700 Subject: [PATCH 1/4] KAFKA-2555: Infinite recursive function call when call commitSync in RebalanceListener.onPartitionRevoked() --- .../apache/kafka/clients/consumer/internals/Coordinator.java | 5 +++-- 1 file changed, 3 insertions(+), 2 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 81a7c9c45fea6..1dfdb6cc6bc65 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 @@ -134,7 +134,6 @@ public void refreshCommittedOffsetsIfNeeded() { public Map fetchCommittedOffsets(Set partitions) { while (true) { ensureCoordinatorKnown(); - ensurePartitionAssignment(); // contact coordinator to fetch committed offsets RequestFuture> future = sendOffsetFetchRequest(partitions); @@ -368,9 +367,11 @@ public void onFailure(RuntimeException e) { } public void commitOffsetsSync(Map offsets) { + if (offsets.isEmpty()) + return; + while (true) { ensureCoordinatorKnown(); - ensurePartitionAssignment(); RequestFuture future = sendOffsetCommitRequest(offsets); client.poll(future); From 72fb2a8d3037cf3fa4bf26690a484bf7a724a32e Mon Sep 17 00:00:00 2001 From: Jiangjie Qin Date: Wed, 23 Sep 2015 20:49:27 -0700 Subject: [PATCH 2/4] Addressed Jason's comments. --- .../apache/kafka/common/errors/IllegalGenerationException.java | 2 +- .../apache/kafka/common/errors/UnknownConsumerIdException.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/errors/IllegalGenerationException.java b/clients/src/main/java/org/apache/kafka/common/errors/IllegalGenerationException.java index d20b74a46fecf..fe8ba7a4bd38b 100644 --- a/clients/src/main/java/org/apache/kafka/common/errors/IllegalGenerationException.java +++ b/clients/src/main/java/org/apache/kafka/common/errors/IllegalGenerationException.java @@ -12,7 +12,7 @@ */ package org.apache.kafka.common.errors; -public class IllegalGenerationException extends RetriableException { +public class IllegalGenerationException extends ApiException { private static final long serialVersionUID = 1L; public IllegalGenerationException() { diff --git a/clients/src/main/java/org/apache/kafka/common/errors/UnknownConsumerIdException.java b/clients/src/main/java/org/apache/kafka/common/errors/UnknownConsumerIdException.java index 9bcbd114e74e8..28bfd72fad8c5 100644 --- a/clients/src/main/java/org/apache/kafka/common/errors/UnknownConsumerIdException.java +++ b/clients/src/main/java/org/apache/kafka/common/errors/UnknownConsumerIdException.java @@ -12,7 +12,7 @@ */ package org.apache.kafka.common.errors; -public class UnknownConsumerIdException extends RetriableException { +public class UnknownConsumerIdException extends ApiException { private static final long serialVersionUID = 1L; public UnknownConsumerIdException() { From 7d905c2d7f93eddae3ce89812a65271623d7d6ed Mon Sep 17 00:00:00 2001 From: Jiangjie Qin Date: Thu, 24 Sep 2015 16:08:14 -0700 Subject: [PATCH 3/4] Addressed Jason's comments. --- .../kafka/clients/consumer/internals/Coordinator.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) 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 1dfdb6cc6bc65..d0d7f393f3481 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 @@ -21,6 +21,8 @@ import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.DisconnectException; +import org.apache.kafka.common.errors.IllegalGenerationException; +import org.apache.kafka.common.errors.UnknownConsumerIdException; import org.apache.kafka.common.metrics.Measurable; import org.apache.kafka.common.metrics.MetricConfig; import org.apache.kafka.common.metrics.Metrics; @@ -196,7 +198,10 @@ private void reassignPartitions() { client.poll(future); if (future.failed()) { - if (!future.isRetriable()) + if (future.exception() instanceof UnknownConsumerIdException + || future.exception() instanceof IllegalGenerationException) + continue; + else if (!future.isRetriable()) throw future.exception(); Utils.sleep(retryBackoffMs); } From 988adc2ba3d393fcf6ba3335004b6f47dd056fac Mon Sep 17 00:00:00 2001 From: Jiangjie Qin Date: Thu, 24 Sep 2015 17:09:58 -0700 Subject: [PATCH 4/4] removed IllegalGenerationIdException from reassignPartitions() --- .../apache/kafka/clients/consumer/internals/Coordinator.java | 3 +-- 1 file changed, 1 insertion(+), 2 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 d0d7f393f3481..0b31c2275aa44 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 @@ -198,8 +198,7 @@ private void reassignPartitions() { client.poll(future); if (future.failed()) { - if (future.exception() instanceof UnknownConsumerIdException - || future.exception() instanceof IllegalGenerationException) + if (future.exception() instanceof UnknownConsumerIdException) continue; else if (!future.isRetriable()) throw future.exception();