Update optime/leader with checkpoint location after clean shut down (#1527)

Potentially this information could be used in order to make sure that there is no data loss on switchover.
This commit is contained in:
Alexander Kukushkin
2020-05-15 16:13:16 +02:00
committed by GitHub
parent 285bffc68d
commit 7cf0b753ab
10 changed files with 52 additions and 16 deletions
+11 -2
View File
@@ -758,10 +758,19 @@ class AbstractDCS(object):
otherwise it should return `!False`""" otherwise it should return `!False`"""
@abc.abstractmethod @abc.abstractmethod
def delete_leader(self): def _delete_leader(self):
"""Voluntarily remove leader key from DCS """Remove leader key from DCS.
This method should remove leader key if current instance is the leader""" This method should remove leader key if current instance is the leader"""
def delete_leader(self, last_operation=None):
"""Update optime/leader and voluntarily remove leader key from DCS.
This method should remove leader key if current instance is the leader.
:param last_operation: latest checkpoint location in bytes"""
if last_operation:
self.write_leader_optime(last_operation)
return self._delete_leader()
@abc.abstractmethod @abc.abstractmethod
def cancel_initialization(self): def cancel_initialization(self):
""" Removes the initialize key for a cluster """ """ Removes the initialize key for a cluster """
+1 -1
View File
@@ -503,7 +503,7 @@ class Consul(AbstractDCS):
return self._client.kv.put(self.history_path, value) return self._client.kv.put(self.history_path, value)
@catch_consul_errors @catch_consul_errors
def delete_leader(self): def _delete_leader(self):
cluster = self.cluster cluster = self.cluster
if cluster and isinstance(cluster.leader, Leader) and cluster.leader.name == self._name: if cluster and isinstance(cluster.leader, Leader) and cluster.leader.name == self._name:
return self._client.kv.delete(self.leader_path, cas=cluster.leader.index) return self._client.kv.delete(self.leader_path, cas=cluster.leader.index)
+1 -1
View File
@@ -625,7 +625,7 @@ class Etcd(AbstractDCS):
return self.retry(self._client.write, self.initialize_path, sysid, prevExist=(not create_new)) return self.retry(self._client.write, self.initialize_path, sysid, prevExist=(not create_new))
@catch_etcd_errors @catch_etcd_errors
def delete_leader(self): def _delete_leader(self):
return self._client.delete(self.leader_path, prevValue=self._name) return self._client.delete(self.leader_path, prevValue=self._name)
@catch_etcd_errors @catch_etcd_errors
+8 -2
View File
@@ -546,9 +546,15 @@ class Kubernetes(AbstractDCS):
resource_version = cluster.config.index if cluster and cluster.config and cluster.config.index else None resource_version = cluster.config.index if cluster and cluster.config and cluster.config.index else None
return self.patch_or_create_config({self._INITIALIZE: sysid}, resource_version) return self.patch_or_create_config({self._INITIALIZE: sysid}, resource_version)
def delete_leader(self): def _delete_leader(self):
"""Unused"""
def delete_leader(self, last_operation=None):
if self.cluster and isinstance(self.cluster.leader, Leader) and self.cluster.leader.name == self._name: if self.cluster and isinstance(self.cluster.leader, Leader) and self.cluster.leader.name == self._name:
self.patch_or_create(self.leader_path, {self._LEADER: None}, self._leader_resource_version, True, False, []) annotations = {self._LEADER: None}
if last_operation:
annotations[self._OPTIME] = last_operation
self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, True, False, [])
self.reset_cluster() self.reset_cluster()
def cancel_initialization(self): def cancel_initialization(self):
+1 -1
View File
@@ -315,7 +315,7 @@ class ZooKeeper(AbstractDCS):
def _update_leader(self): def _update_leader(self):
return True return True
def delete_leader(self): def _delete_leader(self):
self._client.restart() self._client.restart()
return True return True
+8 -6
View File
@@ -748,13 +748,13 @@ class Ha(object):
return self._is_healthiest_node(members.values()) return self._is_healthiest_node(members.values())
def _delete_leader(self): def _delete_leader(self, last_operation=None):
self.set_is_leader(False) self.set_is_leader(False)
self.dcs.delete_leader() self.dcs.delete_leader(last_operation)
self.dcs.reset_cluster() self.dcs.reset_cluster()
def release_leader_key_voluntarily(self): def release_leader_key_voluntarily(self, last_operation=None):
self._delete_leader() self._delete_leader(last_operation)
self.touch_member() self.touch_member()
logger.info("Leader key released") logger.info("Leader key released")
@@ -784,8 +784,9 @@ class Ha(object):
self.set_is_leader(False) self.set_is_leader(False)
if mode_control['release']: if mode_control['release']:
checkpoint_location = self.state_handler.latest_checkpoint_location() if mode == 'graceful' else None
with self._async_executor: with self._async_executor:
self.release_leader_key_voluntarily() self.release_leader_key_voluntarily(checkpoint_location)
time.sleep(2) # Give a time to somebody to take the leader lock time.sleep(2) # Give a time to somebody to take the leader lock
if mode_control['offline']: if mode_control['offline']:
node_to_follow, leader = None, None node_to_follow, leader = None, None
@@ -1384,7 +1385,8 @@ class Ha(object):
stop_timeout=self.master_stop_timeout())) stop_timeout=self.master_stop_timeout()))
if not self.state_handler.is_running(): if not self.state_handler.is_running():
if self.has_lock(): if self.has_lock():
self.dcs.delete_leader() checkpoint_location = self.state_handler.latest_checkpoint_location()
self.dcs.delete_leader(checkpoint_location)
self.touch_member() self.touch_member()
else: else:
# XXX: what about when Patroni is started as the wrong user that has access to the watchdog device # XXX: what about when Patroni is started as the wrong user that has access to the watchdog device
+12 -1
View File
@@ -13,7 +13,7 @@ from patroni.postgresql.bootstrap import Bootstrap
from patroni.postgresql.cancellable import CancellableSubprocess from patroni.postgresql.cancellable import CancellableSubprocess
from patroni.postgresql.config import ConfigHandler from patroni.postgresql.config import ConfigHandler
from patroni.postgresql.connection import Connection, get_connection_cursor from patroni.postgresql.connection import Connection, get_connection_cursor
from patroni.postgresql.misc import parse_history, postgres_major_version_to_int from patroni.postgresql.misc import parse_history, parse_lsn, postgres_major_version_to_int
from patroni.postgresql.postmaster import PostmasterProcess from patroni.postgresql.postmaster import PostmasterProcess
from patroni.postgresql.slots import SlotsHandler from patroni.postgresql.slots import SlotsHandler
from patroni.exceptions import PostgresConnectionException from patroni.exceptions import PostgresConnectionException
@@ -310,6 +310,17 @@ class Postgresql(object):
except (TypeError, ValueError): except (TypeError, ValueError):
logger.exception('Failed to parse timeline from pg_controldata output') logger.exception('Failed to parse timeline from pg_controldata output')
def latest_checkpoint_location(self):
"""Returns checkpoint location for the cleanly shut down primary"""
data = self.controldata()
lsn = data.get('Latest checkpoint location')
if data.get('Database cluster state') == 'shut down' and lsn:
try:
return str(parse_lsn(lsn))
except (IndexError, ValueError) as e:
logger.error('Exception when parsing lsn %s: %r', lsn, e)
def is_running(self): def is_running(self):
"""Returns PostmasterProcess if one is running on the data directory or None. If most recently seen process """Returns PostmasterProcess if one is running on the data directory or None. If most recently seen process
is running updates the cached process based on pid file.""" is running updates the cached process based on pid file."""
+4 -1
View File
@@ -154,7 +154,10 @@ def run_async(self, func, args=()):
@patch.object(Postgresql, '_cluster_info_state_get', Mock(return_value=3)) @patch.object(Postgresql, '_cluster_info_state_get', Mock(return_value=3))
@patch.object(Postgresql, 'call_nowait', Mock(return_value=True)) @patch.object(Postgresql, 'call_nowait', Mock(return_value=True))
@patch.object(Postgresql, 'data_directory_empty', Mock(return_value=False)) @patch.object(Postgresql, 'data_directory_empty', Mock(return_value=False))
@patch.object(Postgresql, 'controldata', Mock(return_value={'Database system identifier': SYSID})) @patch.object(Postgresql, 'controldata', Mock(return_value={
'Database system identifier': SYSID,
'Database cluster state': 'shut down',
'Latest checkpoint location': '0/12345678'}))
@patch.object(SlotsHandler, 'sync_replication_slots', Mock()) @patch.object(SlotsHandler, 'sync_replication_slots', Mock())
@patch.object(ConfigHandler, 'append_pg_hba', Mock()) @patch.object(ConfigHandler, 'append_pg_hba', Mock())
@patch.object(ConfigHandler, 'write_pgpass', Mock(return_value={})) @patch.object(ConfigHandler, 'write_pgpass', Mock(return_value={}))
+1 -1
View File
@@ -113,7 +113,7 @@ class TestKubernetes(unittest.TestCase):
self.k.initialize() self.k.initialize()
def test_delete_leader(self): def test_delete_leader(self):
self.k.delete_leader() self.k.delete_leader(1)
def test_cancel_initialization(self): def test_cancel_initialization(self):
self.k.cancel_initialization() self.k.cancel_initialization()
+5
View File
@@ -338,6 +338,11 @@ class TestPostgresql(BaseTestPostgresql):
with patch.object(Postgresql, '_query', Mock(side_effect=RetryFailedError(''))): with patch.object(Postgresql, '_query', Mock(side_effect=RetryFailedError(''))):
self.assertRaises(PostgresConnectionException, self.p.is_leader) self.assertRaises(PostgresConnectionException, self.p.is_leader)
@patch.object(Postgresql, 'controldata',
Mock(return_value={'Database cluster state': 'shut down', 'Latest checkpoint location': 'X/678'}))
def test_latest_checkpoint_location(self):
self.assertIsNone(self.p.latest_checkpoint_location())
def test_reload(self): def test_reload(self):
self.assertTrue(self.p.reload()) self.assertTrue(self.p.reload())