Merge pull request #26 from CyberDem0n/master

Reduce amount of writes into etcd
This commit is contained in:
Alexander Kukushkin
2015-06-24 16:39:50 +02:00
8 changed files with 87 additions and 26 deletions
+7 -1
View File
@@ -24,9 +24,15 @@ class Governor:
host, port = config['restapi']['listen'].split(':')
self.api = RestApiServer(self, config['restapi'])
self.next_run = time.time()
self.shutdown_member_ttl = 300
def touch_member(self, ttl=None):
connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string
if self.ha.cluster:
for m in self.ha.cluster.members:
# Do not update member TTL when it is far from being expired
if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl:
return True
return self.etcd.touch_member(self.postgresql.name, connection_string, ttl)
def initialize(self):
@@ -94,7 +100,7 @@ def main():
except KeyboardInterrupt:
pass
finally:
governor.touch_member(300) # schedule member removal
governor.touch_member(governor.shutdown_member_ttl) # schedule member removal
governor.postgresql.stop()
governor.etcd.delete_leader(governor.postgresql.name)
+12 -11
View File
@@ -8,7 +8,7 @@ from collections import namedtuple
from dns.exception import DNSException
from dns import resolver
from helpers.errors import CurrentLeaderError, EtcdError, EtcdConnectionFailed
from helpers.utils import sleep
from helpers.utils import calculate_ttl, sleep
from requests.exceptions import RequestException
if sys.hexversion >= 0x03000000:
@@ -19,24 +19,25 @@ else:
logger = logging.getLogger(__name__)
class Member(namedtuple('Member', 'hostname,conn_url,api_url,ttl')):
class Member(namedtuple('Member', 'name,conn_url,api_url,expiration,ttl')):
@staticmethod
def fromNode(node):
scheme, netloc, path, params, query, fragment = urlparse(node['value'])
conn_url = urlunparse((scheme, netloc, path, params, '', fragment))
api_url = None
for name, value in parse_qsl(query):
if name == 'application_name' and value:
api_url = value
break
return Member(node['key'].split('/')[-1], conn_url, api_url, node.get('ttl', None))
api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0]
expiration = node.get('expiration', None)
ttl = node.get('ttl', None)
return Member(node['key'].split('/')[-1], conn_url, api_url, expiration, ttl)
def real_ttl(self):
return calculate_ttl(self.expiration) or -1
class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')):
def is_unlocked(self):
return not (self.leader and self.leader.hostname)
return not (self.leader and self.leader.name)
class Client:
@@ -267,11 +268,11 @@ class Etcd:
node = self.find_node(response['node'], '/leader')
if node:
for m in members:
if m.hostname == node['value']:
if m.name == node['value']:
leader = m
break
if not leader:
leader = Member(node['value'], None, None, None)
leader = Member(node['value'], None, None, None, None)
return Cluster(initialize, leader, last_leader_operation, members)
elif status_code == 404:
+1 -1
View File
@@ -23,7 +23,7 @@ class Ha:
return self.etcd.update_leader(self.state_handler)
def has_lock(self):
lock_owner = self.cluster.leader and self.cluster.leader.hostname
lock_owner = self.cluster.leader and self.cluster.leader.name
logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name)
return lock_owner == self.state_handler.name
+3 -3
View File
@@ -284,7 +284,7 @@ class Postgresql:
return False
for member in cluster.members:
if member.hostname == self.name:
if member.name == self.name:
continue
try:
r = parseurl(member.conn_url)
@@ -297,7 +297,7 @@ class Postgresql:
row = member_cursor.fetchone()
member_cursor.close()
member_conn.close()
logger.error([self.name, member.hostname, row])
logger.error([self.name, member.name, row])
if not row[0] or row[1] < 0:
return False
except psycopg2.Error:
@@ -402,7 +402,7 @@ primary_conninfo = '{}'
self.members = [r[0] for r in cursor]
def create_replication_slots(self, cluster):
members = [m.hostname for m in cluster.members if m.hostname != self.name]
members = [m.name for m in cluster.members if m.name != self.name]
# drop unused slots
for slot in set(self.members) - set(members):
self.query("""SELECT pg_drop_replication_slot(%s)
+35
View File
@@ -1,10 +1,45 @@
import datetime
import os
import re
import signal
import sys
import time
received_sigchld = False
_DATE_TIME_RE = re.compile(r'''^
(?P<year>\d{4})\-(?P<month>\d{2})\-(?P<day>\d{2}) # date
T
(?P<hour>\d{2}):(?P<minute>\d{2}):(?P<second>\d{2})\.(?P<microsecond>\d{6}) # time
\d*Z$''', re.X)
def parse_datetime(time_str):
"""
>>> parse_datetime('2015-06-10T12:56:30.552539016Z')
datetime.datetime(2015, 6, 10, 12, 56, 30, 552539)
>>> parse_datetime('2015-06-10 12:56:30.552539016Z')
"""
m = _DATE_TIME_RE.match(time_str)
if not m:
return None
p = dict((n, int(m.group(n))) for n in 'year month day hour minute second microsecond'.split(' '))
return datetime.datetime(**p)
def calculate_ttl(expiration):
"""
>>> calculate_ttl(None)
>>> calculate_ttl('2015-06-10 12:56:30.552539016Z')
"""
if not expiration:
return None
expiration = parse_datetime(expiration)
if not expiration:
return None
now = datetime.datetime.utcnow()
return int((expiration - now).total_seconds())
def lsn_to_bytes(value):
"""
+14 -1
View File
@@ -1,3 +1,4 @@
import datetime
import dns.resolver
import json
import requests
@@ -7,7 +8,7 @@ import unittest
from dns.exception import DNSException
from helpers.errors import EtcdError, CurrentLeaderError, EtcdConnectionFailed
from helpers.etcd import Client, Cluster, Etcd
from helpers.etcd import Client, Cluster, Etcd, Member
class MockResponse:
@@ -98,6 +99,18 @@ def socket_getaddrinfo(*args):
raise socket.error()
class TestMember(unittest.TestCase):
def __init__(self, method_name='runTest'):
super(TestMember, self).__init__(method_name)
def test_real_ttl(self):
now = datetime.datetime.utcnow()
member = Member('a', 'b', 'c', (now + datetime.timedelta(seconds=2)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None)
self.assertLess(member.real_ttl(), 2)
self.assertEquals(Member('a', 'b', 'c', '', None).real_ttl(), -1)
class TestClient(unittest.TestCase):
def __init__(self, method_name='runTest'):
+9
View File
@@ -1,3 +1,4 @@
import datetime
import psycopg2
import requests
import subprocess
@@ -7,6 +8,7 @@ import unittest
import yaml
from governor import Governor, main
from helpers.etcd import Cluster, Member
from test_ha import true, false
from test_postgresql import Postgresql, subprocess_call, psycopg2_connect
from test_etcd import requests_get, requests_put, requests_delete
@@ -67,6 +69,13 @@ class TestGovernor(unittest.TestCase):
return False
return True
def test_touch_member(self):
now = datetime.datetime.utcnow()
member = Member(self.g.postgresql.name, 'b', 'c', (now + datetime.timedelta(
seconds=self.g.shutdown_member_ttl + 10)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None)
self.g.ha.cluster = Cluster(True, member, 0, [member])
self.g.touch_member()
def test_governor_initialize(self):
self.g.postgresql.should_use_s3_to_create_replica = false
self.g.etcd.client._base_uri = 'http://remote'
+6 -9
View File
@@ -122,7 +122,9 @@ class TestPostgresql(unittest.TestCase):
psycopg2.connect = psycopg2_connect
if not os.path.exists(self.p.data_dir):
os.makedirs(self.p.data_dir)
self.leader = Member('leader', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
self.leader = Member('leader', 'postgres://replicator:[email protected]:5435/postgres', None, None, 28)
self.other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, None, 28)
self.me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, None, 28)
def tear_down(self):
shutil.rmtree('data')
@@ -149,13 +151,11 @@ class TestPostgresql(unittest.TestCase):
self.p.follow_the_leader(None)
self.p.demote(self.leader)
self.p.follow_the_leader(self.leader)
self.p.follow_the_leader(Member('leader', 'postgres://replicator:[email protected]:5435/postgres', None, 28))
self.p.follow_the_leader(self.other)
def test_create_replication_slots(self):
self.p.start()
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, 28)
cluster = Cluster(True, self.leader, 0, [me, other, self.leader])
cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leader])
self.p.create_replication_slots(cluster)
def test_query(self):
@@ -169,10 +169,7 @@ class TestPostgresql(unittest.TestCase):
self.assertRaises(psycopg2.OperationalError, self.p.query, 'blabla')
def test_is_healthiest_node(self):
leader = Member('leader', 'postgres://replicator:[email protected]:5435/postgres', None, 28)
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, 28)
cluster = Cluster(True, leader, 0, [me, other, leader])
cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leader])
self.assertTrue(self.p.is_healthiest_node(cluster))
self.p.is_leader = false
self.assertFalse(self.p.is_healthiest_node(cluster))