From d941f6bc5e3daa34d884804616b342d7234ff5e4 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 20 Dec 2019 12:03:25 +0100 Subject: [PATCH] Use restore_command from standby_cluster config on cascading replicas (#1341) The standby_leader was already doing it from the beginning feature existed. Not doing the same on replicas might prevent them from catching up with standby leader due to WALs being recycled. In addition to that apply the same strategy to archive_cleanup_command. --- patroni/dcs/__init__.py | 10 +++++++--- patroni/ha.py | 16 ++++++++++++++-- patroni/postgresql/config.py | 15 ++++++++------- 3 files changed, 29 insertions(+), 12 deletions(-) diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index a0364d7f..de394b22 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -240,21 +240,25 @@ class Leader(namedtuple('Leader', 'index,session,member')): def conn_url(self): return self.member.conn_url + @property + def data(self): + return self.member.data + @property def timeline(self): - return self.member.data.get('timeline') + return self.data.get('timeline') @property def checkpoint_after_promote(self): """ >>> Leader(1, '', Member.from_node(1, '', '', '{"version":"z"}')).checkpoint_after_promote """ - version = self.member.data.get('version') + version = self.data.get('version') if version: try: # 1.5.6 is the last version which doesn't expose checkpoint_after_promote: false if tuple(map(int, version.split('.'))) > (1, 5, 6): - return self.member.data['role'] == 'master' and 'checkpoint_after_promote' not in self.member.data + return self.data['role'] == 'master' and 'checkpoint_after_promote' not in self.data except Exception: logger.debug('Failed to parse Patroni version %s', version) diff --git a/patroni/ha.py b/patroni/ha.py index 95b76b41..ca012825 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -342,14 +342,26 @@ class Ha(object): def _get_node_to_follow(self, cluster): # determine the node to follow. If replicatefrom tag is set, # try to follow the node mentioned there, otherwise, follow the leader. - if self.is_standby_cluster() and (self.cluster.is_unlocked() or self.has_lock(False)): + standby_config = self.get_standby_cluster_config() + is_standby_cluster = _is_standby_cluster(standby_config) + if is_standby_cluster and (self.cluster.is_unlocked() or self.has_lock(False)): node_to_follow = self.get_remote_master() elif self.patroni.replicatefrom and self.patroni.replicatefrom != self.state_handler.name: node_to_follow = cluster.get_member(self.patroni.replicatefrom) else: node_to_follow = cluster.leader - return node_to_follow if node_to_follow and node_to_follow.name != self.state_handler.name else None + node_to_follow = node_to_follow if node_to_follow and node_to_follow.name != self.state_handler.name else None + + if node_to_follow and not isinstance(node_to_follow, RemoteMember): + # we are going to abuse Member.data to pass following parameters + params = ('restore_command', 'archive_cleanup_command') + for param in params: # It is highly unlikely to happen, but we want to protect from the case + node_to_follow.data.pop(param, None) # when above-mentioned params came from outside. + if is_standby_cluster: + node_to_follow.data.update({p: standby_config[p] for p in params if standby_config.get(p)}) + + return node_to_follow def follow(self, demote_reason, follow_reason, refresh=True): if refresh: diff --git a/patroni/postgresql/config.py b/patroni/postgresql/config.py index 70eafcca..10df64b4 100644 --- a/patroni/postgresql/config.py +++ b/patroni/postgresql/config.py @@ -535,15 +535,16 @@ class ConfigHandler(object): is_remote_master = isinstance(member, RemoteMember) primary_conninfo = self.primary_conninfo_params(member) if primary_conninfo: + use_slots = self.get('use_slots', True) and self._postgresql.major_version >= 90400 + if use_slots and not (is_remote_master and member.no_replication_slot): + primary_slot_name = member.primary_slot_name if is_remote_master else self._postgresql.name + recovery_params['primary_slot_name'] = slot_name_from_member_name(primary_slot_name) recovery_params['primary_conninfo'] = primary_conninfo - if self.get('use_slots', True) and self._postgresql.major_version >= 90400 \ - and not (is_remote_master and member.no_replication_slot): - recovery_params['primary_slot_name'] = member.primary_slot_name if is_remote_master \ - else slot_name_from_member_name(self._postgresql.name) - if is_remote_master: # standby_cluster config might have different parameters, we want to override them - recovery_params.update({p: member.data.get(p) for p in ('restore_command', 'recovery_min_apply_delay', - 'archive_cleanup_command') if member.data.get(p)}) + # standby_cluster config might have different parameters, we want to override them + standby_cluster_params = ['restore_command', 'archive_cleanup_command']\ + + (['recovery_min_apply_delay'] if is_remote_master else []) + recovery_params.update({p: member.data.get(p) for p in standby_cluster_params if member and member.data.get(p)}) return recovery_params def recovery_conf_exists(self):