summaryrefslogtreecommitdiff
path: root/kafka/coordinator/consumer.py
diff options
context:
space:
mode:
Diffstat (limited to 'kafka/coordinator/consumer.py')
-rw-r--r--kafka/coordinator/consumer.py6
1 files changed, 5 insertions, 1 deletions
diff --git a/kafka/coordinator/consumer.py b/kafka/coordinator/consumer.py
index 9d6f4eb..9b7a3cd 100644
--- a/kafka/coordinator/consumer.py
+++ b/kafka/coordinator/consumer.py
@@ -225,7 +225,11 @@ class ConsumerCoordinator(BaseCoordinator):
self._subscription.needs_fetch_committed_offsets = True
# update partition assignment
- self._subscription.assign_from_subscribed(assignment.partitions())
+ try:
+ self._subscription.assign_from_subscribed(assignment.partitions())
+ except ValueError as e:
+ log.warning("%s. Probably due to a deleted topic. Requesting Re-join" % e)
+ self.request_rejoin()
# give the assignor a chance to update internal state
# based on the received assignment