diff --git a/patroni/ha.py b/patroni/ha.py index fc24e32a..f5ebaefe 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -125,7 +125,7 @@ class Ha(object): data['pending_restart'] = True if not self._async_executor.busy and data['state'] in ['running', 'restarting', 'starting']: try: - data['xlog_location'] = self.state_handler.wal_position(retry=False) + data['xlog_location'] = self.state_handler.wal_position() except Exception: pass if self.patroni.scheduled_restart: @@ -993,6 +993,7 @@ class Ha(object): def _run_cycle(self): dcs_failed = False try: + self.state_handler.reset_cluster_info_state() self.load_cluster_from_dcs() if self.is_paused(): diff --git a/patroni/postgresql.py b/patroni/postgresql.py index b1c65c56..d3feaa02 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -36,6 +36,12 @@ STOP_POLLING_INTERVAL = 1 REWIND_STATUS = type('Enum', (), {'INITIAL': 0, 'CHECK': 1, 'NEED': 2, 'NOT_NEED': 3, 'SUCCESS': 4, 'FAILED': 5}) sync_standby_name_re = re.compile('^[A-Za-z_][A-Za-z_0-9\$]*$') +wal_position_query = ("CASE WHEN pg_is_in_recovery() THEN GREATEST(" + " pg_{0}_{1}_diff(COALESCE(pg_last_{0}_receive_{1}(), '0/0'), '0/0')::bigint," + " pg_{0}_{1}_diff(pg_last_{0}_replay_{1}(), '0/0')::bigint)" + " ELSE pg_{0}_{1}_diff(pg_current_{0}_{1}(), '0/0')::bigint " + "END") + def quote_ident(value): """Very simplified version of quote_ident""" @@ -160,6 +166,8 @@ class Postgresql(object): self._state_entry_timestamp = None + self._cluster_info_state = {} + # Last known running process self._postmaster_proc = None @@ -697,13 +705,28 @@ class Postgresql(object): self.set_state('stopped') return ret + def reset_cluster_info_state(self): + self._cluster_info_state = {} + + def _cluster_info_state_get(self, name): + if not self._cluster_info_state: + stmt = "SELECT pg_is_in_recovery(), " + wal_position_query.format(self.wal_name, self.lsn_name) + + try: + result = self._is_leader_retry(self._query, stmt).fetchone() + self._cluster_info_state = dict(zip(['is_in_recovery', 'wal_position'], result)) + except RetryFailedError as e: # SELECT failed two times + self._cluster_info_state = {'error': str(e)} + if not self.is_starting() and self.pg_isready() == STATE_REJECT: + self.set_state('starting') + + if 'error' in self._cluster_info_state: + raise PostgresConnectionException(self._cluster_info_state['error']) + + return self._cluster_info_state.get(name) + def is_leader(self): - try: - return not self._is_leader_retry(self._query, 'SELECT pg_is_in_recovery()').fetchone()[0] - except RetryFailedError as e: # SELECT pg_is_in_recovery() failed two times - if not self.is_starting() and self.pg_isready() == STATE_REJECT: - self.set_state('starting') - raise PostgresConnectionException(str(e)) + return not self._cluster_info_state_get('is_in_recovery') def is_running(self): """Returns PostmasterProcess if one is running on the data directory or None. If most recently seen process @@ -1375,21 +1398,14 @@ BEGIN END; $$""".format(name, ' '.join(options)), name, password, password) - def wal_position(self, retry=True): - stmt = """SELECT CASE WHEN pg_is_in_recovery() - THEN GREATEST(pg_{0}_{1}_diff(COALESCE(pg_last_{0}_receive_{1}(), '0/0'), - '0/0')::bigint, - pg_{0}_{1}_diff(pg_last_{0}_replay_{1}(), '0/0')::bigint) - ELSE pg_{0}_{1}_diff(pg_current_{0}_{1}(), '0/0')::bigint - END""".format(self.wal_name, self.lsn_name) - + def wal_position(self): # This method could be called from different threads (simultaneously with some other `_query` calls). # If it is called not from main thread we will create a new cursor to execute statement. if current_thread().ident == self.__thread_ident: - return (self.query(stmt) if retry else self._query(stmt)).fetchone()[0] + return self._cluster_info_state_get('wal_position') with self.connection().cursor() as cursor: - cursor.execute(stmt) + cursor.execute('SELECT ' + wal_position_query.format(self.wal_name, self.lsn_name)) return cursor.fetchone()[0] def load_replication_slots(self): diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index dbfac543..4c1e48fd 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -35,8 +35,8 @@ class MockCursor(object): self.results = [('blabla',), ('foobar',)] elif sql.startswith('SELECT CASE WHEN pg_is_in_recovery()'): self.results = [(2,)] - elif sql == 'SELECT pg_is_in_recovery()': - self.results = [(False, )] + elif sql.startswith('SELECT pg_is_in_recovery()'): + self.results = [(False, 2)] elif sql.startswith('WITH replication_info AS ('): replication_info = '[{"application_name":"walreceiver","client_addr":"1.2.3.4",' +\ '"state":"streaming","sync_state":"async","sync_priority":0}]' @@ -478,6 +478,7 @@ class TestPostgresql(unittest.TestCase): @patch.object(Postgresql, 'pg_isready', Mock(return_value=STATE_REJECT)) def test_is_leader(self): self.assertTrue(self.p.is_leader()) + self.p.reset_cluster_info_state() with patch.object(Postgresql, '_query', Mock(side_effect=RetryFailedError(''))): self.assertRaises(PostgresConnectionException, self.p.is_leader)