diff --git a/patroni/__init__.py b/patroni/__init__.py index 37a8e848..38cde421 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -84,6 +84,10 @@ class Patroni(object): nap_time = self.next_run - current_time if nap_time <= 0: self.next_run = current_time + # Release the GIL so we don't starve anyone waiting on async_executor lock + time.sleep(0.001) + # Warn user that Patroni is not keeping up + logger.warning("Loop time exceeded, rescheduling immediately.") elif self.dcs.watch(nap_time): self.next_run = time.time() diff --git a/patroni/ha.py b/patroni/ha.py index 93567eed..233b7f9d 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -571,9 +571,7 @@ class Ha(object): # try to start dead postgres if not self.state_handler.is_healthy(): - msg = self.recover() - if msg is not None: - return msg + return self.recover() try: if self.cluster.is_unlocked(): diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 0eec2daa..dc76be03 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -890,7 +890,7 @@ $$""".format(name, ' '.join(options)), name, password, password) def load_replication_slots(self): if self.use_slots and self._schedule_load_slots: - cursor = self.query("SELECT slot_name FROM pg_replication_slots WHERE slot_type='physical'") + cursor = self._query("SELECT slot_name FROM pg_replication_slots WHERE slot_type='physical'") self._replication_slots = [r[0] for r in cursor] self._schedule_load_slots = False @@ -930,19 +930,20 @@ $$""".format(name, ' '.join(options)), name, password, password) # drop unused slots for slot in set(self._replication_slots) - slots: - self.query("""SELECT pg_drop_replication_slot(%s) - WHERE EXISTS(SELECT 1 FROM pg_replication_slots - WHERE slot_name = %s AND NOT active)""", slot, slot) + self._query("""SELECT pg_drop_replication_slot(%s) + WHERE EXISTS(SELECT 1 FROM pg_replication_slots + WHERE slot_name = %s AND NOT active)""", slot, slot) # create new slots for slot in slots - set(self._replication_slots): - self.query("""SELECT pg_create_physical_replication_slot(%s) - WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots - WHERE slot_name = %s)""", slot, slot) + self._query("""SELECT pg_create_physical_replication_slot(%s) + WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots + WHERE slot_name = %s)""", slot, slot) self._replication_slots = slots - except psycopg2.Error: + except Exception: logger.exception('Exception when changing replication slots') + self._schedule_load_slots = True def last_operation(self): return str(self.xlog_position()) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 95a251a2..f8f15b42 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -315,11 +315,9 @@ class TestPostgresql(unittest.TestCase): def test_sync_replication_slots(self): self.p.start() cluster = Cluster(True, None, self.leader, 0, [self.me, self.other, self.leadermem], None) + with mock.patch('patroni.postgresql.Postgresql._query', Mock(side_effect=psycopg2.OperationalError)): + self.p.sync_replication_slots(cluster) self.p.sync_replication_slots(cluster) - self.p.query = Mock(side_effect=psycopg2.OperationalError) - self.p.schedule_load_slots = True - self.p.sync_replication_slots(cluster) - self.p.schedule_load_slots = False with mock.patch('patroni.postgresql.Postgresql.role', new_callable=PropertyMock(return_value='replica')): self.p.sync_replication_slots(cluster) with mock.patch('patroni.postgresql.logger.error', new_callable=Mock()) as errorlog_mock: