Add private _refresh_on_disconnects flag to KafkaClient
This commit is contained in:
@@ -94,6 +94,7 @@ class KafkaClient(object):
|
|||||||
self._metadata_refresh_in_progress = False
|
self._metadata_refresh_in_progress = False
|
||||||
self._conns = {}
|
self._conns = {}
|
||||||
self._connecting = set()
|
self._connecting = set()
|
||||||
|
self._refresh_on_disconnects = True
|
||||||
self._delayed_tasks = DelayedTaskQueue()
|
self._delayed_tasks = DelayedTaskQueue()
|
||||||
self._last_bootstrap = 0
|
self._last_bootstrap = 0
|
||||||
self._bootstrap_fails = 0
|
self._bootstrap_fails = 0
|
||||||
@@ -164,9 +165,10 @@ class KafkaClient(object):
|
|||||||
|
|
||||||
# Connection failures imply that our metadata is stale, so let's refresh
|
# Connection failures imply that our metadata is stale, so let's refresh
|
||||||
elif conn.state is ConnectionStates.DISCONNECTING:
|
elif conn.state is ConnectionStates.DISCONNECTING:
|
||||||
log.warning("Node %s connect failed -- refreshing metadata", node_id)
|
|
||||||
if node_id in self._connecting:
|
if node_id in self._connecting:
|
||||||
self._connecting.remove(node_id)
|
self._connecting.remove(node_id)
|
||||||
|
if self._refresh_on_disconnects:
|
||||||
|
log.warning("Node %s connect failed -- refreshing metadata", node_id)
|
||||||
self.cluster.request_update()
|
self.cluster.request_update()
|
||||||
|
|
||||||
def _maybe_connect(self, node_id):
|
def _maybe_connect(self, node_id):
|
||||||
@@ -597,9 +599,13 @@ class KafkaClient(object):
|
|||||||
if node_id is None:
|
if node_id is None:
|
||||||
raise Errors.NoBrokersAvailable()
|
raise Errors.NoBrokersAvailable()
|
||||||
|
|
||||||
|
# We will be intentionally causing socket failures
|
||||||
|
# and should not trigger metadata refresh
|
||||||
|
self._refresh_on_disconnects = False
|
||||||
self._maybe_connect(node_id)
|
self._maybe_connect(node_id)
|
||||||
conn = self._conns[node_id]
|
conn = self._conns[node_id]
|
||||||
version = conn.check_version()
|
version = conn.check_version()
|
||||||
|
self._refresh_on_disconnects = True
|
||||||
return version
|
return version
|
||||||
|
|
||||||
def wakeup(self):
|
def wakeup(self):
|
||||||
|
|||||||
Reference in New Issue
Block a user