From 3b0ea900cd92762f39c2f656ae2d20e7abf20efe Mon Sep 17 00:00:00 2001 From: Christopher Winslett Date: Wed, 13 May 2015 16:37:14 -0700 Subject: [PATCH 1/9] followers without a leader based on feedback from @CyberDem0n and closes #3 --- governor.py | 2 +- helpers/ha.py | 1 + helpers/postgresql.py | 19 +++++++++++++++---- 3 files changed, 17 insertions(+), 5 deletions(-) diff --git a/governor.py b/governor.py index 31081f07..c3847bd2 100755 --- a/governor.py +++ b/governor.py @@ -55,7 +55,7 @@ if postgresql.data_directory_empty(): else: time.sleep(5) else: - postgresql.write_recovery_conf({"address": "postgres://169.0.0.1:5432"}) + postgresql.follow_no_leader() postgresql.start() while True: diff --git a/helpers/ha.py b/helpers/ha.py index 80736420..bd770524 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -59,6 +59,7 @@ class Ha: self.state_handler.demote(self.fetch_current_leader()) return "demoting self because i am not the healthiest node" elif self.fetch_current_leader() is None: + self.state_handler.follow_no_leader() return "waiting on leader to be elected because i am not the healthiest node" else: self.state_handler.follow_the_leader(self.fetch_current_leader()) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 387ff19a..f6bbcd69 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -158,15 +158,18 @@ class Postgresql: f.close() def write_recovery_conf(self, leader_hash): - leader = urlparse(leader_hash["address"]) - f = open("%s/recovery.conf" % self.data_dir, "w") f.write(""" standby_mode = 'on' primary_slot_name = '%(recovery_slot)s' -primary_conninfo = 'user=%(user)s password=%(password)s host=%(hostname)s port=%(port)s sslmode=prefer sslcompression=1' recovery_target_timeline = 'latest' -""" % {"recovery_slot": self.name, "user": leader.username, "password": leader.password, "hostname": leader.hostname, "port": leader.port}) +""" % {"recovery_slot": self.name}) + if leader_hash is not None: + leader = urlparse(leader_hash["address"]) + f.write(""" +primary_conninfo = 'user=%(user)s password=%(password)s host=%(hostname)s port=%(port)s sslmode=prefer sslcompression=1' + """ % {"user": leader.username, "password": leader.password, "hostname": leader.hostname, "port": leader.port}) + if "recovery_conf" in self.config: for name, value in self.config["recovery_conf"].iteritems(): f.write("%s = '%s'\n" % (name, value)) @@ -179,6 +182,14 @@ recovery_target_timeline = 'latest' self.restart() return True + def follow_no_leader(self): + print "initing leaderless follower" + if os.system("grep primary_conninfo %(data_dir)s/recovery.conf > /dev/null" % {"data_dir": self.data_dir}) == 0: + self.write_recovery_conf(None) + if self.is_running(): + self.restart() + return True + def promote(self): return os.system("pg_ctl promote -w -D %s" % self.data_dir) == 0 From 96a889f358d82d070690465eae1838322d191e2d Mon Sep 17 00:00:00 2001 From: Christopher Winslett Date: Wed, 13 May 2015 16:56:25 -0700 Subject: [PATCH 2/9] fix logic for dead leaders returning online without checking that recovery.conf exists, the prior logic would not create the recovyer.conf, and a dead leader would return to a primary state --- helpers/postgresql.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index f6bbcd69..fd93f954 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -92,7 +92,8 @@ class Postgresql: logger.info("Removed %s" % pid_path) command_code = os.system("postgres -D %s %s &" % (self.data_dir, self.server_options())) - time.sleep(5) + while not self.is_running(): + time.sleep(5) return command_code != 0 def stop(self): @@ -183,8 +184,7 @@ primary_conninfo = 'user=%(user)s password=%(password)s host=%(hostname)s port=% return True def follow_no_leader(self): - print "initing leaderless follower" - if os.system("grep primary_conninfo %(data_dir)s/recovery.conf > /dev/null" % {"data_dir": self.data_dir}) == 0: + if not os.path.exists("%s/recovery.conf" % self.data_dir) or os.system("grep primary_conninfo %(data_dir)s/recovery.conf &> /dev/null" % {"data_dir": self.data_dir}) == 0: self.write_recovery_conf(None) if self.is_running(): self.restart() From bec1a3c11ee464adb06aae446ccae414c100865b Mon Sep 17 00:00:00 2001 From: Christopher Winslett Date: Thu, 14 May 2015 13:02:07 -0700 Subject: [PATCH 3/9] support synchronous replication --- README.md | 23 +++++++++++++++++++++++ governor.py | 1 - helpers/postgresql.py | 4 ++++ 3 files changed, 27 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 5a0fa978..7a5fa765 100644 --- a/README.md +++ b/README.md @@ -58,6 +58,29 @@ For an example file, see `postgres0.yml`. Below is an explanation of settings: * *recovery_conf*: configuration settings written to recovery.conf when configuring follower * *parameters*: list of configuration settings for Postgres +## Replication choices + +Governor uses Postgres' streaming replication. By default, this replication is asynchronous. For more information, see the (Postgres documentation on streaming replication)[http://www.postgresql.org/docs/current/static/warm-standby.html#STREAMING-REPLICATION]. + +Governor's asynchronous replication configuration allows for `maximum_lag_on_failover` settings. This setting ensures replication will not occur if a follower is more than a certain number of bytes behind the follower. This setting should be increased or decreased based on business requirements. + +When asynchronous replication is not best for your use-case, investigate how Postgres's (synchronous replication)[http://www.postgresql.org/docs/current/static/warm-standby.html#SYNCHRONOUS-REPLICATION] works. Synchronous replication ensures consistency across a cluster by confirming that writes are written to a secondary before returning to the connecting client with a success. The cost of synchronous replication will be reduced throughput on writes. This throughput will be entirely based on network performance. In hosted datacenter environments (like AWS, Rackspace, or any network you do not control), synchrous replication increases the variability of write performance significantly. If followers become inaccessible from the leader, the leader will becomes effectively readonly. + +To enable a simple synchronous replication test, add the follow lines to the `parameters` section of your YAML configuration files. + +```YAML + synchronous_commit: "on" + synchronous_standby_names: "*" +``` + +When using synchronous replication, use at least a 3-Postgres data nodes to ensure write availability if one host fails. + +Choosing your replication schema is dependent on the many business decisions. Investigate both async and sync replication, as well as other HA solutions, to determine which solution is best for you. + +## Applications should not use superusers + +When connecting from an application, always use a non-superuser. Governor requires access to the database to function properly. By using a superuser from application, you can potentially use the entire connection pool, including the connections reserved for superusers with the `superuser_reserved_connections` setting. If Governor cannot access the Primary, because the connection pool is full, behavior will be undesireable. + ## Requirements on a Mac Run the following on a Mac to install requirements: diff --git a/governor.py b/governor.py index c3847bd2..ee90eaf1 100755 --- a/governor.py +++ b/governor.py @@ -40,7 +40,6 @@ if postgresql.data_directory_empty(): postgresql.initialize() etcd.take_leader(postgresql.name) postgresql.start() - postgresql.create_replication_user() else: synced_from_leader = False while not synced_from_leader: diff --git a/helpers/postgresql.py b/helpers/postgresql.py index fd93f954..b30f3c9c 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -56,6 +56,10 @@ class Postgresql: def initialize(self): if os.system("initdb -D %s" % self.data_dir) == 0: + # start Postgres without options to setup replication user indepedent of other system settings + os.system("pg_ctl start -w -D %s" % self.data_dir) + self.create_replication_user() + os.system("pg_ctl stop -w -m fast -D %s" % self.data_dir) self.write_pg_hba() return True From a4e40959f91388ad2306f1ebb495d682c5b28915 Mon Sep 17 00:00:00 2001 From: Christopher Winslett Date: Thu, 14 May 2015 14:18:52 -0700 Subject: [PATCH 4/9] fix linking syntax in README --- README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 7a5fa765..5af483f8 100644 --- a/README.md +++ b/README.md @@ -60,11 +60,11 @@ For an example file, see `postgres0.yml`. Below is an explanation of settings: ## Replication choices -Governor uses Postgres' streaming replication. By default, this replication is asynchronous. For more information, see the (Postgres documentation on streaming replication)[http://www.postgresql.org/docs/current/static/warm-standby.html#STREAMING-REPLICATION]. +Governor uses Postgres' streaming replication. By default, this replication is asynchronous. For more information, see the [Postgres documentation on streaming replication](http://www.postgresql.org/docs/current/static/warm-standby.html#STREAMING-REPLICATION). Governor's asynchronous replication configuration allows for `maximum_lag_on_failover` settings. This setting ensures replication will not occur if a follower is more than a certain number of bytes behind the follower. This setting should be increased or decreased based on business requirements. -When asynchronous replication is not best for your use-case, investigate how Postgres's (synchronous replication)[http://www.postgresql.org/docs/current/static/warm-standby.html#SYNCHRONOUS-REPLICATION] works. Synchronous replication ensures consistency across a cluster by confirming that writes are written to a secondary before returning to the connecting client with a success. The cost of synchronous replication will be reduced throughput on writes. This throughput will be entirely based on network performance. In hosted datacenter environments (like AWS, Rackspace, or any network you do not control), synchrous replication increases the variability of write performance significantly. If followers become inaccessible from the leader, the leader will becomes effectively readonly. +When asynchronous replication is not best for your use-case, investigate how Postgres's [synchronous replication](http://www.postgresql.org/docs/current/static/warm-standby.html#SYNCHRONOUS-REPLICATION) works. Synchronous replication ensures consistency across a cluster by confirming that writes are written to a secondary before returning to the connecting client with a success. The cost of synchronous replication will be reduced throughput on writes. This throughput will be entirely based on network performance. In hosted datacenter environments (like AWS, Rackspace, or any network you do not control), synchrous replication increases the variability of write performance significantly. If followers become inaccessible from the leader, the leader will becomes effectively readonly. To enable a simple synchronous replication test, add the follow lines to the `parameters` section of your YAML configuration files. From d221d1de1c65d0fc02ca789a9c26dea0e0bfaa43 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 18 May 2015 10:53:27 +0200 Subject: [PATCH 5/9] listen_address can have more then one value separated by comma. We will use the first one to connect --- helpers/postgresql.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index f872385d..1c0e29b1 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -30,6 +30,7 @@ class Postgresql: def __init__(self, config): self.name = config['name'] + self.listen_addresses, self.port = config['listen'].split(':') self.data_dir = config['data_dir'] self.replication = config['replication'] self.recovery_conf = os.path.join(self.data_dir, 'recovery.conf') @@ -46,7 +47,8 @@ class Postgresql: def cursor(self): if not self.cursor_holder: - self.conn = psycopg2.connect('postgres://{}/postgres'.format(self.config['listen'])) + self.conn = psycopg2.connect('postgres://{}:{}/postgres'.format( + self.listen_addresses.split(',')[0].strip(), self.port)) self.conn.autocommit = True self.cursor_holder = self.conn.cursor() @@ -131,8 +133,7 @@ class Postgresql: return os.system(self._pg_ctl + ' restart -m fast') == 0 def server_options(self): - host, port = self.config['listen'].split(':') - options = '--listen_addresses={} --port={}'.format(host, port) + options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port) for setting, value in self.config['parameters'].items(): options += " --{}='{}'".format(setting, value) return options From b24fb0488c777bebb4ce97a7e97c529cd852e673 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 18 May 2015 13:33:27 +0200 Subject: [PATCH 6/9] last_leader_operation is the propery of Cluster object and the value is set in get_cluster method --- helpers/etcd.py | 49 +++++++++++++++++++------------------------ helpers/ha.py | 8 +------ helpers/postgresql.py | 17 +++++---------- 3 files changed, 28 insertions(+), 46 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index cbe44507..29552a2b 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -8,14 +8,8 @@ from helpers.errors import CurrentLeaderError, EtcdError logger = logging.getLogger(__name__) -class Member(namedtuple('Member', 'hostname,address')): - - pass - - -class Cluster(namedtuple('Cluster', 'leader,members')): - - pass +Member = namedtuple('Member', 'hostname,address,ttl') +Cluster = namedtuple('Cluster', 'leader,last_leader_operation,members') class Etcd: @@ -87,21 +81,32 @@ class Etcd: 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 [] + # get list of members + node = self.find_node(response['node'], '/members') or {'nodes': []} + members = [Member(n['key'].split('/')[-1], n['value'], n.get('ttl', None)) for n in node['nodes']] - leader_node = self.find_node(response['node'], '/leader') - if leader_node: + # get last leader operation + last_leader_operation = 0 + node = self.find_node(response['node'], '/optime') + if node: + node = self.find_node(node, '/leader') + if node: + last_leader_operation = int(node['value']) + + # get leader + leader = None + node = self.find_node(response['node'], '/leader') + if node: for m in members: - if m.hostname == leader_node['value']: + if m.hostname == node['value']: leader = m break if not leader: - leader = Member(leader['value'], None) - return Cluster(leader, members) + leader = Member(leader['value'], None, None) + + return Cluster(leader, last_leader_operation, members) elif status_code == 404: - return Cluster(None, []) + return Cluster(None, None, []) except: logger.exception('get_cluster') @@ -135,16 +140,6 @@ class Etcd: def race(self, path, value): return self.put_client_path(path, value=value, prevExist=False) - def last_leader_operation(self): - try: - response, status_code = self.get_client_path('/optime/leader') - if status_code == 404: - return None - return int(response['node']['value']) - except: - logger.exception('last_leader_operation') - raise EtcdError('Etcd is not responding properly') - def delete_member(self, member): return self.delete_client_path('/members/' + member) diff --git a/helpers/ha.py b/helpers/ha.py index a2788e4c..89467bf4 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -1,5 +1,4 @@ import logging -import time from helpers.errors import EtcdError, HealthiestMemberError from psycopg2 import OperationalError @@ -50,7 +49,7 @@ class Ha: self.load_cluster_from_etcd() if self.is_unlocked(): - if self.state_handler.is_healthiest_node(self.etcd.last_leader_operation(), self.cluster.members): + if self.state_handler.is_healthiest_node(self.cluster): if self.acquire_lock(): if not self.state_handler.is_leader(): self.state_handler.promote() @@ -97,8 +96,3 @@ class Ha: logger.error('Error communicating with Postgresql. Will try again') except HealthiestMemberError: logger.error('failed to determine healthiest member fromt etcd') - - def run(self): - while True: - self.run_cycle() - time.sleep(10) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 1c0e29b1..7b8a15c8 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -1,7 +1,6 @@ import logging import os import psycopg2 -import re import sys import time @@ -144,11 +143,11 @@ class Postgresql: return False return True - def is_healthiest_node(self, last_leader_operation, members): - if (last_leader_operation or 0) - self.xlog_position() > self.config.get('maximum_lag_on_failover', 0): + def is_healthiest_node(self, cluster): + if cluster.last_leader_operation - self.xlog_position() > self.config.get('maximum_lag_on_failover', 0): return False - for member in members: + for member in cluster.members: if member.hostname == self.name: continue try: @@ -159,19 +158,13 @@ class Postgresql: "SELECT %s - (pg_last_xlog_replay_location() - '0/0000000'::pg_lsn)", (self.xlog_position(), )) xlog_diff = member_cursor.fetchone()[0] logger.info([self.name, member.hostname, xlog_diff]) - if xlog_diff < 0: - member_cursor.close() - return False member_cursor.close() + if xlog_diff < 0: + return False except psycopg2.OperationalError: continue return True - def replication_slot_name(self): - member = os.environ.get("MEMBER") - (member, _) = re.subn(r'[^a-z0-9]+', r'_', member) - return member - def write_pg_hba(self): with open(os.path.join(self.data_dir, 'pg_hba.conf'), 'a') as f: f.write('host replication {username} {network} md5'.format(**self.replication)) From 422512880f23f0d58ded235a69dc8d50b87b4583 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 18 May 2015 13:43:14 +0200 Subject: [PATCH 7/9] last_leader_operation is the propery of Cluster object and the value is set in get_cluster method --- helpers/etcd.py | 1 + helpers/postgresql.py | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index 29552a2b..ecffc9ec 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -92,6 +92,7 @@ class Etcd: node = self.find_node(node, '/leader') if node: last_leader_operation = int(node['value']) + print(last_leader_operation) # get leader leader = None diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 7b8a15c8..4fd16880 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -79,7 +79,7 @@ class Postgresql: return not os.path.exists(self.data_dir) or os.listdir(self.data_dir) == [] def initialize(self): - if os.system(self._pg_ctl + ' initdb') == 0: + if os.system(self._pg_ctl + ' initdb -o --encoding=UTF8') == 0: self.write_pg_hba() return True From 520de12232574a6d859c2fc05f6dd9c58bb8b91a Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 18 May 2015 14:39:14 +0200 Subject: [PATCH 8/9] Remove debug print --- helpers/etcd.py | 1 - 1 file changed, 1 deletion(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index ecffc9ec..29552a2b 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -92,7 +92,6 @@ class Etcd: node = self.find_node(node, '/leader') if node: last_leader_operation = int(node['value']) - print(last_leader_operation) # get leader leader = None From f6fc60800a9561f17f5bd929407b917e68c96542 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Mon, 18 May 2015 16:54:36 +0200 Subject: [PATCH 9/9] Close connection when querying other members of cluster. --- helpers/postgresql.py | 1 + 1 file changed, 1 insertion(+) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index d5a86cc7..ecbbcdbd 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -167,6 +167,7 @@ class Postgresql: xlog_diff = member_cursor.fetchone()[0] logger.info([self.name, member.hostname, xlog_diff]) member_cursor.close() + member_conn.close() if xlog_diff < 0: return False except psycopg2.OperationalError: