mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 07:30:14 +00:00
Merge branch 'master' of github.com:zalando/patroni into feature/quorum-commit
This commit is contained in:
+99
-96
@@ -204,6 +204,20 @@ class Ha(object):
|
||||
if not value:
|
||||
self._promote_timestamp = 0
|
||||
|
||||
def sync_mode_is_active(self) -> bool:
|
||||
"""Check whether synchronous replication is requested and already active.
|
||||
|
||||
:returns: ``True`` if the primary already put its name into the ``/sync`` in DCS.
|
||||
"""
|
||||
return self.is_synchronous_mode() and not self.cluster.sync.is_empty
|
||||
|
||||
def quorum_commit_mode_is_active(self) -> bool:
|
||||
"""Checks whether quorum replication is requested and already active.
|
||||
|
||||
:returns: ``True`` if the primary already put its name into the ``/sync`` in DCS.
|
||||
"""
|
||||
return self.is_quorum_commit_mode() and not self.cluster.sync.is_empty
|
||||
|
||||
def load_cluster_from_dcs(self) -> None:
|
||||
cluster = self.dcs.get_cluster()
|
||||
|
||||
@@ -469,7 +483,7 @@ class Ha(object):
|
||||
if timeout == 0:
|
||||
# We are requested to prefer failing over to restarting primary. But see first if there
|
||||
# is anyone to fail over to.
|
||||
if self.is_failover_possible(self.cluster.members):
|
||||
if self.is_failover_possible():
|
||||
self.watchdog.disable()
|
||||
logger.info("Primary crashed. Failing over.")
|
||||
self.demote('immediate')
|
||||
@@ -624,18 +638,10 @@ class Ha(object):
|
||||
""":returns: `True` if synchronous replication is requested."""
|
||||
return self.global_config.is_synchronous_mode
|
||||
|
||||
def is_synchronous_mode_active(self) -> bool:
|
||||
""":returns: `True` is synchronous replication requested is active (/sync key has a valid "leader" field)."""
|
||||
return self.is_synchronous_mode() and not self.cluster.sync.is_empty
|
||||
|
||||
def is_quorum_commit_mode(self) -> bool:
|
||||
""":returns: `True` if quorum commit replication is requested and "supported"."""
|
||||
return self.global_config.is_quorum_commit_mode and self.state_handler.supports_multiple_sync
|
||||
|
||||
def is_quorum_commit_mode_active(self) -> bool:
|
||||
""":returns: `True` if quorum replication is requested and active (/sync key has a valid "leader" field)."""
|
||||
return self.is_quorum_commit_mode() and not self.cluster.sync.is_empty
|
||||
|
||||
def is_failsafe_mode(self) -> bool:
|
||||
""":returns: `True` if failsafe_mode is enabled in global configuration."""
|
||||
return self.global_config.check_mode('failsafe_mode')
|
||||
@@ -789,7 +795,7 @@ class Ha(object):
|
||||
self.disable_synchronous_replication()
|
||||
return True
|
||||
|
||||
if self.is_quorum_commit_mode_active():
|
||||
if self.quorum_commit_mode_is_active():
|
||||
sync = CaseInsensitiveSet(self.cluster.sync.members)
|
||||
numsync = len(sync) - self.cluster.sync.quorum - 1
|
||||
if self.state_handler.name not in sync: # Node outside voters achieved quorum and got leader
|
||||
@@ -951,6 +957,8 @@ class Ha(object):
|
||||
return _MemberStatus.unknown(member)
|
||||
|
||||
def fetch_nodes_statuses(self, members: List[Member]) -> List[_MemberStatus]:
|
||||
if not members:
|
||||
return []
|
||||
pool = ThreadPool(len(members))
|
||||
results = pool.map(self.fetch_node_status, members) # Run API calls on members in parallel
|
||||
pool.close()
|
||||
@@ -1034,7 +1042,7 @@ class Ha(object):
|
||||
logger.info('My timeline %s is behind last known cluster timeline %s', my_timeline, cluster_timeline)
|
||||
return False
|
||||
|
||||
if self.is_quorum_commit_mode_active():
|
||||
if self.quorum_commit_mode_is_active():
|
||||
quorum = self.cluster.sync.quorum
|
||||
voting_set = CaseInsensitiveSet(self.cluster.sync.members)
|
||||
else:
|
||||
@@ -1055,63 +1063,60 @@ class Ha(object):
|
||||
quorum_votes = 0 if self.state_handler.name in voting_set else -1
|
||||
nodes_ahead = 0
|
||||
|
||||
if members:
|
||||
for st in self.fetch_nodes_statuses(members):
|
||||
if st.failover_limitation() is None:
|
||||
if st.in_recovery is False:
|
||||
logger.warning('Primary (%s) is still alive', st.member.name)
|
||||
for st in self.fetch_nodes_statuses(members):
|
||||
if st.failover_limitation() is None:
|
||||
if st.in_recovery is False:
|
||||
logger.warning('Primary (%s) is still alive', st.member.name)
|
||||
return False
|
||||
if my_wal_position < st.wal_position:
|
||||
nodes_ahead += 1
|
||||
logger.info('Wal position of %s is ahead of my wal position', st.member.name)
|
||||
# In synchronous mode the former leader might be still accessible and even be ahead of us.
|
||||
# We should not disqualify himself from the leader race in such a situation.
|
||||
if not self.sync_mode_is_active() or not self.cluster.sync.leader_matches(st.member.name):
|
||||
return False
|
||||
if my_wal_position < st.wal_position:
|
||||
nodes_ahead += 1
|
||||
logger.info('Wal position of %s is ahead of my wal position', st.member.name)
|
||||
# In synchronous mode the former leader might be still accessible and even be ahead of us.
|
||||
# We should not disqualify himself from the leader race in such a situation.
|
||||
if not self.is_synchronous_mode_active()\
|
||||
or not self.cluster.sync.leader_matches(st.member.name):
|
||||
return False
|
||||
logger.info('Ignoring the former leader being ahead of us')
|
||||
# we want to count votes only from nodes with postgres up and running!
|
||||
elif st.member.name in voting_set and st.wal_position > 0:
|
||||
logger.info('Got quorum vote from %s', st.member.name)
|
||||
quorum_votes += 1
|
||||
logger.info('Ignoring the former leader being ahead of us')
|
||||
# we want to count votes only from nodes with postgres up and running!
|
||||
elif st.member.name in voting_set and st.wal_position > 0:
|
||||
logger.info('Got quorum vote from %s', st.member.name)
|
||||
quorum_votes += 1
|
||||
|
||||
# When not in quorum commit we just want to return `True`.
|
||||
# In quorum commit the former leader is special and counted healthy even when there are no other nodes.
|
||||
# Otherwise check that the number of votes exceeds the quorum field from the /sync key.
|
||||
return not self.is_quorum_commit_mode_active() or quorum_votes >= quorum\
|
||||
return not self.quorum_commit_mode_is_active() or quorum_votes >= quorum\
|
||||
or nodes_ahead == 0 and self.cluster.sync.leader == self.state_handler.name
|
||||
|
||||
def is_failover_possible(self, members: List[Member], check_synchronous: Optional[bool] = True,
|
||||
cluster_lsn: Optional[int] = 0) -> bool:
|
||||
"""Checks whether one of the members from the list can possibly win the leader race.
|
||||
def is_failover_possible(self, *, cluster_lsn: int = 0, exclude_failover_candidate: bool = False) -> bool:
|
||||
"""Checks whether any of the cluster members is allowed to promote and is healthy enough for that.
|
||||
|
||||
:param members: list of members to check
|
||||
:param check_synchronous: consider only members that are known to be listed in /sync key when sync replication.
|
||||
:param cluster_lsn: to calculate replication lag and exclude member if it is laggin
|
||||
:returns: `True` if there are members eligible to be the new leader
|
||||
:param cluster_lsn: to calculate replication lag and exclude member if it is lagging.
|
||||
:param exclude_failover_candidate: if ``True``, exclude :attr:`failover.candidate` from the members
|
||||
list against which the failover possibility checks are run.
|
||||
:returns: `True` if there are members eligible to become the new leader.
|
||||
"""
|
||||
candidates = self.get_failover_candidates(exclude_failover_candidate)
|
||||
|
||||
if self.is_synchronous_mode() and self.cluster.failover and self.cluster.failover.candidate and not candidates:
|
||||
logger.warning('Failover candidate=%s does not match with sync_standbys=%s',
|
||||
self.cluster.failover.candidate, self.cluster.sync.sync_standby)
|
||||
elif not candidates:
|
||||
logger.warning('manual failover: candidates list is empty')
|
||||
|
||||
ret = False
|
||||
cluster_timeline = self.cluster.timeline
|
||||
members = [m for m in members if m.name != self.state_handler.name and not m.nofailover and m.api_url]
|
||||
if check_synchronous and self.is_synchronous_mode_active():
|
||||
# If quorum commit is requested we want to check all nodes (even not voters),
|
||||
# because they could get enough votes and reach necessary quorum + 1.
|
||||
members = [m for m in members if self.is_quorum_commit_mode() or self.cluster.sync.matches(m.name)]
|
||||
if members:
|
||||
for st in self.fetch_nodes_statuses(members):
|
||||
not_allowed_reason = st.failover_limitation()
|
||||
if not_allowed_reason:
|
||||
logger.info('Member %s is %s', st.member.name, not_allowed_reason)
|
||||
elif cluster_lsn and st.wal_position < cluster_lsn or\
|
||||
not cluster_lsn and self.is_lagging(st.wal_position):
|
||||
logger.info('Member %s exceeds maximum replication lag', st.member.name)
|
||||
elif self.check_timeline() and (not st.timeline or st.timeline < cluster_timeline):
|
||||
logger.info('Timeline %s of member %s is behind the cluster timeline %s',
|
||||
st.timeline, st.member.name, cluster_timeline)
|
||||
else:
|
||||
ret = True
|
||||
else:
|
||||
logger.warning('manual failover: members list is empty')
|
||||
for st in self.fetch_nodes_statuses(candidates):
|
||||
not_allowed_reason = st.failover_limitation()
|
||||
if not_allowed_reason:
|
||||
logger.info('Member %s is %s', st.member.name, not_allowed_reason)
|
||||
elif cluster_lsn and st.wal_position < cluster_lsn or \
|
||||
not cluster_lsn and self.is_lagging(st.wal_position):
|
||||
logger.info('Member %s exceeds maximum replication lag', st.member.name)
|
||||
elif self.check_timeline() and (not st.timeline or st.timeline < cluster_timeline):
|
||||
logger.info('Timeline %s of member %s is behind the cluster timeline %s',
|
||||
st.timeline, st.member.name, cluster_timeline)
|
||||
else:
|
||||
ret = True
|
||||
return ret
|
||||
|
||||
def manual_failover_process_no_leader(self) -> Optional[bool]:
|
||||
@@ -1161,9 +1166,8 @@ class Ha(object):
|
||||
# try to pick some other members to failover and check that they are healthy
|
||||
if failover.leader:
|
||||
if self.state_handler.name == failover.leader: # I was the leader
|
||||
# exclude me (leader) and desired member which is unhealthy (failover.candidate can be None)
|
||||
members = [m for m in self.cluster.members if m.name not in (failover.candidate, failover.leader)]
|
||||
if self.is_failover_possible(members): # check that there are healthy members
|
||||
# exclude desired member which is unhealthy if it was specified
|
||||
if self.is_failover_possible(exclude_failover_candidate=bool(failover.candidate)):
|
||||
return False
|
||||
else: # I was the leader and it looks like currently I am the only healthy member
|
||||
return True
|
||||
@@ -1216,8 +1220,8 @@ class Ha(object):
|
||||
|
||||
if self.cluster.failover:
|
||||
# When doing a switchover in synchronous mode only synchronous nodes and former leader are allowed to race
|
||||
if self.is_synchronous_mode() and self.cluster.failover.leader and \
|
||||
not self.cluster.sync.is_empty and not self.cluster.sync.matches(self.state_handler.name, True):
|
||||
if self.sync_mode_is_active() and not self.cluster.sync.matches(self.state_handler.name, True) and \
|
||||
self.cluster.failover.leader:
|
||||
return False
|
||||
return self.manual_failover_process_no_leader() or False
|
||||
|
||||
@@ -1238,7 +1242,7 @@ class Ha(object):
|
||||
all_known_members += self.cluster.members
|
||||
|
||||
# Special handling if synchronous mode was requested and activated (the leader in /sync is not empty)
|
||||
if self.is_synchronous_mode_active():
|
||||
if self.sync_mode_is_active():
|
||||
# In quorum commit mode we allow nodes outside of "voters" to take part in
|
||||
# the leader race. They just need to get enough votes to `reach quorum + 1`.
|
||||
if not self.is_quorum_commit_mode() and not self.cluster.sync.matches(self.state_handler.name, True):
|
||||
@@ -1291,9 +1295,7 @@ class Ha(object):
|
||||
# It could happen if Postgres is still archiving the backlog of WAL files.
|
||||
# If we know that there are replicas that received the shutdown checkpoint
|
||||
# location, we can remove the leader key and allow them to start leader race.
|
||||
|
||||
# for a manual failover/switchover with a candidate, we should check the requested candidate only
|
||||
if self.is_failover_possible(self.get_failover_candidates(), cluster_lsn=checkpoint_location):
|
||||
if self.is_failover_possible(cluster_lsn=checkpoint_location):
|
||||
self.state_handler.set_role('demoted')
|
||||
with self._async_executor:
|
||||
self.release_leader_key_voluntarily(checkpoint_location)
|
||||
@@ -1395,19 +1397,11 @@ class Ha(object):
|
||||
if not failover.candidate or failover.candidate != self.state_handler.name:
|
||||
if not failover.candidate and self.is_paused():
|
||||
logger.warning('Failover is possible only to a specific candidate in a paused state')
|
||||
elif self.is_failover_possible():
|
||||
ret = self._async_executor.try_run_async('manual failover: demote', self.demote, ('graceful',))
|
||||
return ret or 'manual failover: demoting myself'
|
||||
else:
|
||||
if self.is_synchronous_mode():
|
||||
members = self.get_failover_candidates(check_sync=True)
|
||||
if failover.candidate and not members:
|
||||
logger.warning('Failover candidate=%s does not match with sync_standbys=%s',
|
||||
failover.candidate, self.cluster.sync.sync_standby)
|
||||
else:
|
||||
members = self.get_failover_candidates()
|
||||
if self.is_failover_possible(members, False): # check that there are healthy members
|
||||
ret = self._async_executor.try_run_async('manual failover: demote', self.demote, ('graceful',))
|
||||
return ret or 'manual failover: demoting myself'
|
||||
else:
|
||||
logger.warning('manual failover: no healthy members found, failover is not possible')
|
||||
logger.warning('manual failover: no healthy members found, failover is not possible')
|
||||
else:
|
||||
logger.warning('manual failover: I am already the leader, no need to failover')
|
||||
else:
|
||||
@@ -1668,7 +1662,7 @@ class Ha(object):
|
||||
if self.has_lock() and self.update_lock():
|
||||
if self._async_executor.scheduled_action == 'doing crash recovery in a single user mode':
|
||||
time_left = self.global_config.primary_start_timeout - (time.time() - self._crash_recovery_started)
|
||||
if time_left <= 0 and self.is_failover_possible(self.cluster.members):
|
||||
if time_left <= 0 and self.is_failover_possible():
|
||||
logger.info("Demoting self because crash recovery is taking too long")
|
||||
self.state_handler.cancellable.cancel(True)
|
||||
self.demote('immediate')
|
||||
@@ -1780,7 +1774,7 @@ class Ha(object):
|
||||
time_left = timeout - self.state_handler.time_in_state()
|
||||
|
||||
if time_left <= 0:
|
||||
if self.is_failover_possible(self.cluster.members):
|
||||
if self.is_failover_possible():
|
||||
logger.info("Demoting self because primary startup is taking too long")
|
||||
self.demote('immediate')
|
||||
return 'stopped PostgreSQL because of startup timeout'
|
||||
@@ -2058,8 +2052,7 @@ class Ha(object):
|
||||
# If we know that there are replicas that received the shutdown checkpoint
|
||||
# location, we can remove the leader key and allow them to start leader race.
|
||||
|
||||
# for a manual failover/switchover with a candidate, we should check the requested candidate only
|
||||
if self.is_failover_possible(self.get_failover_candidates(), cluster_lsn=checkpoint_location):
|
||||
if self.is_failover_possible(cluster_lsn=checkpoint_location):
|
||||
self.dcs.delete_leader(checkpoint_location)
|
||||
status['deleted'] = True
|
||||
else:
|
||||
@@ -2120,23 +2113,33 @@ class Ha(object):
|
||||
name = member.name if member else 'remote_member:{}'.format(uuid.uuid1())
|
||||
return RemoteMember(name, data)
|
||||
|
||||
def get_failover_candidates(self, check_sync: bool = False) -> List[Member]:
|
||||
"""Return list of candidates for either manual or automatic failover.
|
||||
def get_failover_candidates(self, exclude_failover_candidate: bool) -> List[Member]:
|
||||
"""Return a list of candidates for either manual or automatic failover.
|
||||
|
||||
Mainly used to later be passed to ``Ha.is_failover_possible()``.
|
||||
Exclude non-sync members when in synchronous mode, the current node (its checks are always performed earlier)
|
||||
and the candidate if required. If failover candidate exclusion is not requested and a candidate is specified
|
||||
in the /failover key, return the candidate only.
|
||||
The result is further evaluated in the caller :func:`Ha.is_failover_possible` to check if any member is actually
|
||||
healthy enough and is allowed to poromote.
|
||||
|
||||
:param check_sync: if ``True``, also check against the sync key members
|
||||
:param exclude_failover_candidate: if ``True``, exclude :attr:`failover.candidate` from the candidates.
|
||||
|
||||
:returns: a list of ``Member`` ojects or an empty list if there is no candidate available
|
||||
:returns: a list of :class:`Member` ojects or an empty list if there is no candidate available.
|
||||
"""
|
||||
failover = self.cluster.failover
|
||||
if check_sync:
|
||||
exclude = [self.state_handler.name] + ([failover.candidate] if failover and exclude_failover_candidate else [])
|
||||
|
||||
def is_eligible(node: Member) -> bool:
|
||||
# If quorum commit is requested we want to check all nodes (even not voters),
|
||||
# because they could get enough votes and reach necessary quorum + 1.
|
||||
# TODO: allow manual failover (=no leader specified) to async node
|
||||
# every sync_standby or the candidate specified if is in sync_standbys
|
||||
return [m for m in self.cluster.members
|
||||
if self.cluster.sync.matches(m.name)
|
||||
and (not failover or not failover.candidate or m.name == failover.candidate)]
|
||||
else:
|
||||
# every member or the candidate specified
|
||||
return [m for m in self.cluster.members
|
||||
if not failover or not failover.candidate or m.name == failover.candidate]
|
||||
if self.sync_mode_is_active() \
|
||||
and not (self.is_quorum_commit_mode() or self.cluster.sync.matches(node.name)):
|
||||
return False
|
||||
# Don't spend time on "nofailover" nodes checking.
|
||||
# We also don't need nodes which we can't query with the api in the list.
|
||||
return node.name not in exclude and \
|
||||
not node.nofailover and bool(node.api_url) and \
|
||||
(not failover or not failover.candidate or node.name == failover.candidate)
|
||||
|
||||
return list(filter(is_eligible, self.cluster.members))
|
||||
|
||||
@@ -750,7 +750,6 @@ class TestHa(PostgresInit):
|
||||
self.p.is_leader = true
|
||||
self.ha.has_lock = true
|
||||
self.ha.is_synchronous_mode = true
|
||||
self.ha.is_failover_possible = false
|
||||
self.ha.process_sync_replication = Mock()
|
||||
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, self.p.name, 'a', None), (self.p.name, None))
|
||||
self.assertEqual('no action. I am (postgresql0), the leader with the lock', self.ha.run_cycle())
|
||||
|
||||
Reference in New Issue
Block a user