make healthiest node work

This commit is contained in:
Christopher Winslett
2015-03-16 09:47:52 -07:00
parent 450b9912d7
commit 0d75152d54
3 changed files with 35 additions and 8 deletions
+14 -6
View File
@@ -37,12 +37,6 @@ class Etcd:
def client_url(self, path):
return "http://%s/v2/keys/service/%s%s" % (self.host, self.scope, path)
def xlog_position(member):
try:
return self.get_client_path("/service/postgresql/xlog-position/%s" % member)["node"]["value"]
except urllib2.HTTPError:
return None
def current_leader(self):
try:
hostname = self.get_client_path("/leader")["node"]["value"]
@@ -54,6 +48,20 @@ class Etcd:
return None
raise helpers.errors.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:
return None
raise helpers.errors.CurrentLeaderError("Etcd is not responding properly")
def touch_member(self, member, connection_string):
self.put_client_path("/members/%s" % member, {"value": connection_string})
+1 -1
View File
@@ -33,7 +33,7 @@ class Ha:
try:
if self.state_handler.is_healthy():
if self.is_unlocked():
if self.state_handler.is_healthiest_node():
if self.state_handler.is_healthiest_node(self.etcd.members()):
if self.acquire_lock():
if not self.state_handler.is_leader():
self.state_handler.promote()
+20 -1
View File
@@ -104,7 +104,23 @@ class Postgresql:
return True
def is_healthiest_node(self):
def is_healthiest_node(self, members):
for member in members:
if member["hostname"] == self.name:
continue
try:
member_conn = psycopg2.connect(member["address"])
member_conn.autocommit = True
member_cursor = member_conn.cursor()
member_cursor.execute("SELECT '%s'::pg_lsn - pg_last_xlog_replay_location() AS bytes;" % self.xlog_position())
xlog_diff = member_cursor.fetchone()[0]
print [self.name, member["hostname"], xlog_diff]
if xlog_diff < 0:
member_cursor.close()
return False
member_cursor.close()
except psycopg2.OperationalError:
continue
return True
def replication_slot_name(self):
@@ -150,3 +166,6 @@ recovery_target_timeline = 'latest'
def create_replication_user(self):
self.query("CREATE USER \"%s\" WITH REPLICATION ENCRYPTED PASSWORD '%s';" % (self.replication["username"], self.replication["password"]))
def xlog_position(self):
return self.query("SELECT pg_last_xlog_replay_location();").fetchone()[0]