diff options
| author | Dana Powers <dana.powers@gmail.com> | 2016-02-18 21:48:14 -0800 |
|---|---|---|
| committer | Dana Powers <dana.powers@gmail.com> | 2016-02-18 21:48:14 -0800 |
| commit | 6946aa29106eaea4db6dc0166909be590db9d276 (patch) | |
| tree | ca305646b0ba43550f8d3e93b775c2f095cbb2e3 | |
| parent | 3147dfd64493c12c519104bf4751e00871b2c619 (diff) | |
| download | kafka-python-6946aa29106eaea4db6dc0166909be590db9d276.tar.gz | |
Verify node ready before sending offset fetch request from coordinator
| -rw-r--r-- | kafka/coordinator/consumer.py | 5 |
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) |
