From f494d2ce646dfdc35771eb8472dfb543e264829e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 14 Sep 2015 11:19:46 +0200 Subject: [PATCH 1/4] Build Cluster object for ZooKeeper the same way as for Etcd Previous implementation was always setting Cluster.initialize to True. Also it was throwing ZooKeeperError when there were no members in a cluster. Plus BUGFIX of a bug introduced with https://github.com/zalando/patroni/pull/34 in a `load_members` method. - data = self.get_node(self.member_path) + data = self.get_node(self.members_path + member) It was always fetching the same node for all cluster members. Fortunately Etcd doesn't have such problem because we are fetching the whole cluster directory with one recursive API call. --- patroni/zookeeper.py | 50 ++++++++++++++++++++++++++--------------- tests/test_zookeeper.py | 10 +++++++++ 2 files changed, 42 insertions(+), 18 deletions(-) diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 7c01f278..5f9625ac 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -92,9 +92,8 @@ class ZooKeeper(AbstractDCS): self.client.add_listener(self.session_listener) self.cluster_event = self.client.handler.event_object() + self.cluster = None self.fetch_cluster = True - self.members = [] - self.leader = None self.last_leader_operation = 0 self.client.start(None) @@ -121,18 +120,35 @@ class ZooKeeper(AbstractDCS): conn_url, api_url = parse_connection_string(value) return Member(znode.mzxid, name, conn_url, api_url, None, None) + def get_children(self, key, watch=None): + try: + return self.client.get_children(key, watch) + except NoNodeError: + pass + except: + logger.exception('get_children') + return [] + def load_members(self): members = [] - for member in self.client.get_children(self.members_path, self.cluster_watcher): - data = self.get_node(self.member_path) + for member in self.get_children(self.members_path, self.cluster_watcher): + data = self.get_node(self.members_path + member) if data is not None: members.append(self.member(member, *data)) return members def _inner_load_cluster(self): self.cluster_event.clear() - leader = self.get_node(self.leader_path, self.cluster_watcher) - self.members = self.load_members() + nodes = set(self.get_children(self.client_path(''))) + + # get initialize flag + initialize = self._INITIALIZE in nodes + + # get list of members + members = self.load_members() if self._MEMBERS[:-1] in nodes else [] + + # get leader + leader = self.get_node(self.leader_path, self.cluster_watcher) if self._LEADER in nodes else None if leader: client_id = self.client.client_id if leader[0] == self._name and client_id is not None and client_id[0] != leader[1].ephemeralOwner: @@ -142,15 +158,14 @@ class ZooKeeper(AbstractDCS): if leader: member = Member(-1, leader[0], None, None, None, None) - member = ([m for m in self.members if m.name == leader[0]] or [member])[0] + member = ([m for m in members if m.name == leader[0]] or [member])[0] leader = Leader(leader[1].mzxid, None, None, member) self.fetch_cluster = member.index == -1 - self.leader = leader - if self.fetch_cluster: - last_leader_operation = self.get_node(self.leader_optime_path) - if last_leader_operation: - self.last_leader_operation = int(last_leader_operation[0]) + # get last leader operation + self.last_leader_operation = self.get_node(self.leader_optime_path) if self.fetch_cluster else None + self.last_leader_operation = 0 if self.last_leader_operation is None else int(self.last_leader_operation[0]) + self.cluster = Cluster(initialize, leader, self.last_leader_operation, members) def get_cluster(self): if self.exhibitor and self.exhibitor.poll(): @@ -163,7 +178,7 @@ class ZooKeeper(AbstractDCS): logger.exception('get_cluster') self.session_listener(KazooState.LOST) raise ZooKeeperError('ZooKeeper in not responding properly') - return Cluster(True, self.leader, self.last_leader_operation, self.members) + return self.cluster def _create(self, path, value, **kwargs): try: @@ -181,9 +196,8 @@ class ZooKeeper(AbstractDCS): return self._create(self.initialize_path, self._name, makepath=True) def touch_member(self, connection_string, ttl=None): - for m in self.members: - if m.name == self._name: - return True + if self.cluster and any(m.name == self._name for m in self.cluster.members): + return True path = self.member_path try: self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True) @@ -217,8 +231,8 @@ class ZooKeeper(AbstractDCS): return True def delete_leader(self): - if isinstance(self.leader, Leader) and self.leader.name == self._name: - self.client.delete(self.leader_path) + if isinstance(self.cluster, Cluster) and self.cluster.leader.name == self._name: + self.client.delete(self.leader_path, version=self.cluster.leader.index) def cancel_initialization(self): node = self.get_node(self.initialize_path) diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 0a95d0bc..19435c6e 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -75,6 +75,12 @@ class MockKazooClient: return ('foo', ZnodeStat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0)) def get_children(self, path, watch=None, include_data=False): + if path == '/no_node': + raise NoNodeError + elif path == '/other_exception': + raise Exception() + elif path in ['/service/bla/', '/service/test/']: + return ['initialize', 'leader', 'members', 'optime'] return ['foo', 'bar', 'buzz'] def create(self, path, value="", acl=None, ephemeral=False, sequence=False, makepath=False): @@ -138,6 +144,10 @@ class TestZooKeeper(unittest.TestCase): self.assertIsNone(self.zk.get_node('/no_node')) self.assertIsNone(self.zk.get_node('/other_exception')) + def test_get_children(self): + self.assertListEqual(self.zk.get_children('/no_node'), []) + self.assertListEqual(self.zk.get_children('/other_exception'), []) + def test__inner_load_cluster(self): self.zk._base_path = self.zk._base_path.replace('test', 'bla') self.zk._inner_load_cluster() From 209c985420b5b437b8afc9dfedfabf7366a01944 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 14 Sep 2015 11:45:00 +0200 Subject: [PATCH 2/4] get_node and get_children should catch only NoNodeError exception. All other exceptions are needed to have retry functionality working correctly. --- patroni/zookeeper.py | 10 ++-------- tests/test_zookeeper.py | 6 ------ 2 files changed, 2 insertions(+), 14 deletions(-) diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 5f9625ac..f2ef2e99 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -110,10 +110,7 @@ class ZooKeeper(AbstractDCS): try: return self.client.get(key, watch) except NoNodeError: - pass - except: - logger.exception('get_node') - return None + return None @staticmethod def member(name, value, znode): @@ -124,10 +121,7 @@ class ZooKeeper(AbstractDCS): try: return self.client.get_children(key, watch) except NoNodeError: - pass - except: - logger.exception('get_children') - return [] + return [] def load_members(self): members = [] diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 19435c6e..4b31bad2 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -58,8 +58,6 @@ class MockKazooClient: def get(self, path, watch=None): if path == '/no_node': raise NoNodeError - elif path == '/other_exception': - raise Exception() elif '/members/' in path: return ( 'postgres://repuser:rep-pass@localhost:5434/postgres?application_name=http://127.0.0.1:8009/patroni', @@ -77,8 +75,6 @@ class MockKazooClient: def get_children(self, path, watch=None, include_data=False): if path == '/no_node': raise NoNodeError - elif path == '/other_exception': - raise Exception() elif path in ['/service/bla/', '/service/test/']: return ['initialize', 'leader', 'members', 'optime'] return ['foo', 'bar', 'buzz'] @@ -142,11 +138,9 @@ class TestZooKeeper(unittest.TestCase): def test_get_node(self): self.assertIsNone(self.zk.get_node('/no_node')) - self.assertIsNone(self.zk.get_node('/other_exception')) def test_get_children(self): self.assertListEqual(self.zk.get_children('/no_node'), []) - self.assertListEqual(self.zk.get_children('/other_exception'), []) def test__inner_load_cluster(self): self.zk._base_path = self.zk._base_path.replace('test', 'bla') From 4a081bcb7179ad00f3771113176c27706ca43a76 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 14 Sep 2015 11:58:10 +0200 Subject: [PATCH 3/4] Run cancel_initialization with retry --- patroni/zookeeper.py | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index f2ef2e99..38891c58 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -228,13 +228,16 @@ class ZooKeeper(AbstractDCS): if isinstance(self.cluster, Cluster) and self.cluster.leader.name == self._name: self.client.delete(self.leader_path, version=self.cluster.leader.index) - def cancel_initialization(self): + def _cancel_initialization(self): node = self.get_node(self.initialize_path) if node and node[0] == self._name: - try: - self.client.retry(self.client.delete, self.initialize_path, version=node[1].mzxid) - except KazooException: - logger.exception("Unable to delete initialize key") + self.client.delete(self.initialize_path, version=node[1].mzxid) + + def cancel_initialization(self): + try: + self.client.retry(self._cancel_initialization) + except: + logger.exception("Unable to delete initialize key") def watch(self, timeout): self.cluster_event.wait(timeout) From 98488a00a232d0294a2c4f5edc6414bf5a271d2b Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 14 Sep 2015 12:00:24 +0200 Subject: [PATCH 4/4] Remove unused import of KazooException --- patroni/zookeeper.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 38891c58..708cbafa 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -4,7 +4,7 @@ import requests import time from kazoo.client import KazooClient, KazooState -from kazoo.exceptions import NoNodeError, NodeExistsError, KazooException +from kazoo.exceptions import NoNodeError, NodeExistsError from patroni.dcs import AbstractDCS, Cluster, DCSError, Leader, Member, parse_connection_string from patroni.utils import sleep from requests.exceptions import RequestException