summaryrefslogtreecommitdiff
path: root/kafka/coordinator
diff options
context:
space:
mode:
authorDana Powers <dana.powers@gmail.com>2016-02-18 21:48:14 -0800
committerDana Powers <dana.powers@gmail.com>2016-02-18 21:48:14 -0800
commit6946aa29106eaea4db6dc0166909be590db9d276 (patch)
treeca305646b0ba43550f8d3e93b775c2f095cbb2e3 /kafka/coordinator
parent3147dfd64493c12c519104bf4751e00871b2c619 (diff)
downloadkafka-python-6946aa29106eaea4db6dc0166909be590db9d276.tar.gz
Verify node ready before sending offset fetch request from coordinator
Diffstat (limited to 'kafka/coordinator')
-rw-r--r--kafka/coordinator/consumer.py5
1 files changed, 5 insertions, 0 deletions
diff --git a/kafka/coordinator/consumer.py b/kafka/coordinator/consumer.py
index d63d052..b3ff56d 100644
--- a/kafka/coordinator/consumer.py
+++ b/kafka/coordinator/consumer.py
@@ -561,6 +561,11 @@ class ConsumerCoordinator(BaseCoordinator):
else:
node_id = self._client.least_loaded_node()
+ # Verify node is ready
+ if not self._client.ready(node_id):
+ log.debug("Node %s not ready -- failing offset fetch request")
+ return Future().failure(Errors.NodeNotReadyError)
+
log.debug("Fetching committed offsets for partitions: %s", partitions)
# construct the request
topic_partitions = collections.defaultdict(set)