Merge pull request #37 from zalando/feature/zookeeper-fetch-initialize

Build Cluster object for ZooKeeper the same way as for Etcd
This commit is contained in:
Oleksii Kliukin
2015-09-14 12:55:10 +02:00
2 changed files with 46 additions and 31 deletions
+39 -28
View File
@@ -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
@@ -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)
@@ -111,28 +110,39 @@ 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):
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:
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 +152,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 +172,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 +190,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,16 +225,19 @@ 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):
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)
+7 -3
View File
@@ -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',
@@ -75,6 +73,10 @@ 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 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):
@@ -136,7 +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'), [])
def test__inner_load_cluster(self):
self.zk._base_path = self.zk._base_path.replace('test', 'bla')