From eb1721f65ff51b803910705a677eade1d425628c Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 11 May 2015 15:15:08 +0200 Subject: [PATCH] Minimize amount of requests to etcd --- governor.py | 5 -- helpers/etcd.py | 163 +++++++++++++++++++++--------------------- helpers/ha.py | 55 ++++++++------ helpers/postgresql.py | 20 +++--- 4 files changed, 128 insertions(+), 115 deletions(-) diff --git a/governor.py b/governor.py index e145531c..bbca37a0 100755 --- a/governor.py +++ b/governor.py @@ -67,9 +67,4 @@ else: while True: logging.info(ha.run_cycle()) - # create replication slots - if postgresql.is_leader(): - members = [m['hostname'] for m in etcd.members() if m['hostname'] != postgresql.name] - postgresql.create_replication_slots(members) - time.sleep(config["loop_wait"]) diff --git a/helpers/etcd.py b/helpers/etcd.py index e3603b0c..7e8182f8 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -1,117 +1,120 @@ -import urllib2 -import json -import time import logging +import requests +import time -from helpers.errors import CurrentLeaderError -from urllib import urlencode +from collections import namedtuple +from helpers.errors import CurrentLeaderError, EtcdError logger = logging.getLogger(__name__) +class Member(namedtuple('Member', 'hostname,address')): + + pass + + +class Cluster(namedtuple('Cluster', 'leader,members')): + + pass + + class Etcd: def __init__(self, config): self.ttl = config['ttl'] self.base_client_url = 'http://{host}/v2/keys/service/{scope}'.format(**config) + self.postgres_cluster = None def get_client_path(self, path, max_attempts=1): attempts = 0 response = None while True: + ex = None try: - response = urllib2.urlopen(self.client_url(path)).read() - break - except (urllib2.HTTPError, urllib2.URLError) as e: - attempts += 1 - if attempts < max_attempts: - logger.info('Failed to return %s, trying again. (%s of %s)', path, attempts, max_attempts) - time.sleep(3) - else: - raise e - try: - return json.loads(response) - except ValueError: - return response + response = requests.get(self.client_url(path)) + if response.status_code == 200: + break + except Exception, e: + logger.exception('get_client_path') + ex = e - def put_client_path(self, path, data): - opener = urllib2.build_opener(urllib2.HTTPHandler) - request = urllib2.Request(self.client_url(path), data=urlencode(data).replace("false", "False")) - request.get_method = lambda: 'PUT' - opener.open(request) + attempts += 1 + if attempts < max_attempts: + logger.info('Failed to return %s, trying again. (%s of %s)', path, attempts, max_attempts) + time.sleep(3) + elif ex: + raise ex + + return response.json(), response.status_code + + def put_client_path(self, path, **data): + try: + response = requests.put(self.client_url(path), data=data) + return response.status_code in [200, 201] + except: + logger.exception('PUT %s data=%s', path, data) + return False def client_url(self, path): return self.base_client_url + path + @staticmethod + def find_node(node, key): + if not node['dir']: + return None + key = node['key'] + key + for n in node['nodes']: + if n['key'] == key: + return n + return None + + def get_cluster(self): + try: + response, status_code = self.get_client_path('?recursive=true') + if status_code == 200: + leader = None + members = self.find_node(response['node'], '/members') + members = [Member(n['key'].split('/')[-1], n['value']) for n in members['nodes']] if members else [] + + leader_node = self.find_node(response['node'], '/leader') + if leader_node: + for m in members: + if m.hostname == leader_node['value']: + leader = m + break + if not leader: + leader = Member(leader['value'], None) + return Cluster(leader, members) + elif status_code == 404: + return Cluster(None, []) + except: + logger.exception('get_cluster') + + raise EtcdError('Etcd is not responding properly') + def current_leader(self): try: - hostname = self.get_client_path('/leader')['node']['value'] - address = self.get_client_path('/members/' + hostname)['node']['value'] - - return {'hostname': hostname, 'address': address} - except urllib2.HTTPError as e: - if e.code == 404: - return None - raise CurrentLeaderError("Etcd is not responding properly") - - def members(self): - try: - members = [] - - r = self.get_client_path("/members?recursive=true") - for node in r["node"]["nodes"]: - members.append({"hostname": node["key"].split('/')[-1], "address": node["value"]}) - - return members - except urllib2.HTTPError as e: - if e.code == 404: + cluster = self.get_cluster() + if not cluster['leader'] or not cluster['leader'].address: return None + return cluster['leader'] + except: raise CurrentLeaderError("Etcd is not responding properly") def touch_member(self, member, connection_string): - self.put_client_path('/members/' + member, {"value": connection_string}) + self.put_client_path('/members/' + member, value=connection_string) def take_leader(self, value): - return self.put_client_path("/leader", {"value": value, "ttl": self.ttl}) is None + return self.put_client_path('/leader', value=value, ttl=self.ttl) def attempt_to_acquire_leader(self, value): - try: - return self.put_client_path("/leader", {"value": value, "ttl": self.ttl, "prevExist": False}) is None - except urllib2.HTTPError as e: - if e.code == 412: - logger.info('Could not take out TTL lock: %s', e) - return False + ret = self.put_client_path('/leader', value=value, ttl=self.ttl, prevExist=False) + ret or logger.info('Could not take out TTL lock') + return ret def update_leader(self, value): - try: - self.put_client_path("/leader", {"value": value, "ttl": self.ttl, "prevValue": value}) - return True - except urllib2.HTTPError: - logger.error("Error updating TTL on ETCD for primary.") - return False - - def leader_unlocked(self): - try: - self.get_client_path("/leader") - return False - except urllib2.HTTPError as e: - if e.code == 404: - return True - return False - except ValueError as e: - return False - - def am_i_leader(self, value): - # try: - reponse = self.get_client_path("/leader") - logger.info('Lock owner: %s; I am %s', reponse["node"]["value"], value) - return reponse["node"]["value"] == value - # except Exception as e: - # return False + return self.put_client_path('/leader', value=value, ttl=self.ttl, prevValue=value) def race(self, path, value): - try: - return self.put_client_path(path, {"prevExist": False, "value": value}) is None - except urllib2.HTTPError: - return False + return self.put_client_path(path, value=value, prevExist=False) diff --git a/helpers/ha.py b/helpers/ha.py index 633ececf..da8e0d18 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -2,7 +2,7 @@ import inspect import logging import time -from helpers.errors import CurrentLeaderError, HealthiestMemberError +from helpers.errors import EtcdError, HealthiestMemberError from psycopg2 import OperationalError logger = logging.getLogger(__name__) @@ -18,6 +18,10 @@ class Ha: def __init__(self, state_handler, etcd): self.state_handler = state_handler self.etcd = etcd + self.cluster = None + + def load_cluster_from_etcd(self): + self.cluster = self.etcd.get_cluster() def acquire_lock(self): return self.etcd.attempt_to_acquire_leader(self.state_handler.name) @@ -26,60 +30,71 @@ class Ha: return self.etcd.update_leader(self.state_handler.name) def is_unlocked(self): - return self.etcd.leader_unlocked() + return not (self.cluster.leader and self.cluster.leader.hostname) def has_lock(self): - return self.etcd.am_i_leader(self.state_handler.name) + logger.info('Lock owner: %s; I am %s', self.cluster.leader.hostname, self.state_handler.name) + return self.cluster.leader.hostname == self.state_handler.name - def fetch_current_leader(self): - return self.etcd.current_leader() + def demote(self): + return self.state_handler.demote(self.cluster.leader) + + def follow_the_leader(self): + return self.state_handler.follow_the_leader(self.cluster.leader) def run_cycle(self): try: if self.state_handler.is_healthy(): + self.load_cluster_from_etcd() if self.is_unlocked(): - if self.state_handler.is_healthiest_node(self.etcd.members()): + if self.state_handler.is_healthiest_node(self.cluster.members): if self.acquire_lock(): if not self.state_handler.is_leader(): self.state_handler.promote() return "promoted self to leader by acquiring session lock" return "acquired session lock as a leader" else: + self.load_cluster_from_etcd() if self.state_handler.is_leader(): - self.state_handler.demote(self.fetch_current_leader()) + self.demote() return "demoted self due after trying and failing to obtain lock" else: - self.state_handler.follow_the_leader(self.fetch_current_leader()) + self.follow_the_leader() return "following new leader after trying and failing to obtain lock" else: + self.load_cluster_from_etcd() if self.state_handler.is_leader(): - self.state_handler.demote(self.fetch_current_leader()) + self.demote() return "demoting self because i am not the healthiest node" else: - self.state_handler.follow_the_leader(self.fetch_current_leader()) + self.follow_the_leader() return "following a different leader because i am not the healthiest node" else: if self.has_lock() and self.update_lock(): - if not self.state_handler.is_leader(): - self.state_handler.promote() - return "promoted self to leader because i had the session lock" - else: - return "no action. i am the leader with the lock" + try: + if not self.state_handler.is_leader(): + self.state_handler.promote() + return "promoted self to leader because i had the session lock" + else: + return "no action. i am the leader with the lock" + finally: + # create replication slots + self.state_handler.create_replication_slots([m.hostname for m in self.cluster.members]) else: logger.info("does not have lock") if self.state_handler.is_leader(): - self.state_handler.demote(self.fetch_current_leader()) + self.demote() return "demoting self because i do not have the lock and i was a leader" else: - self.state_handler.follow_the_leader(self.fetch_current_leader()) + self.follow_the_leader() return "no action. i am a secondary and i am following a leader" else: - if not self.state_handler.is_running(): # XXX is_running == is_healthy + if not self.state_handler.is_running(): # XXX is_running == is_healthy self.state_handler.start() return "postgresql was stopped. starting again." return "no action. not healthy enough to do anything." - except CurrentLeaderError: - logger.error("failed to fetch current leader from etcd") + except EtcdError: + logger.error("Error communicating with Etcd") except OperationalError: logger.error("Error communicating with Postgresql. Will try again.") except HealthiestMemberError: diff --git a/helpers/postgresql.py b/helpers/postgresql.py index dd8ad4c6..a4a368af 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -139,10 +139,10 @@ class Postgresql: def is_healthiest_node(self, members): for member in members: - if member['hostname'] == self.name: + if member.hostname == self.name: continue try: - member_conn = psycopg2.connect(member['address']) + member_conn = psycopg2.connect(member.address) member_conn.autocommit = True member_cursor = member_conn.cursor() member_cursor.execute( @@ -171,11 +171,11 @@ class Postgresql: r = parseurl(leader_url) return 'user={username} password={password} host={hostname} port={port} sslmode=prefer sslcompression=1'.format(**r) - def check_recovery_conf(self, leader_hash): + def check_recovery_conf(self, leader): if not os.path.isfile(self.recovery_conf): return False - pattern = leader_hash and 'address' in leader_hash and self.primary_conninfo(leader_hash['address']) + pattern = leader and leader.address and self.primary_conninfo(leader.address) with open(self.recovery_conf, 'r') as f: for line in f: @@ -186,23 +186,23 @@ class Postgresql: return not pattern - def write_recovery_conf(self, leader_hash): + def write_recovery_conf(self, leader): with open(self.recovery_conf, 'w') as f: f.write("""standby_mode = 'on' recovery_target_timeline = 'latest' """) - if leader_hash and 'address' in leader_hash: + if leader and leader.address: f.write(""" primary_slot_name = '{}' primary_conninfo = '{}' -""".format(self.name, self.primary_conninfo(leader_hash['address']))) +""".format(self.name, self.primary_conninfo(leader.address))) for name, value in self.config.get('recovery_conf', {}).iteritems(): f.write("{} = '{}'\n".format(name, value)) - def follow_the_leader(self, leader_hash): - if self.check_recovery_conf(leader_hash): + def follow_the_leader(self, leader): + if self.check_recovery_conf(leader): return - self.write_recovery_conf(leader_hash) + self.write_recovery_conf(leader) self.restart() def promote(self):