From dcdd71384b729f03f71c6becee558e777163cbfb Mon Sep 17 00:00:00 2001 From: chenxiaobin <1025221611@qq.com> Date: Fri, 16 Apr 2021 14:38:36 +0800 Subject: [PATCH] optimize primary hang scenario when standby in catchup --- .../storage/replication/syncrep.cpp | 123 ++++++++---------- .../storage/replication/walsender.cpp | 48 ++++++- src/include/replication/walsender_private.h | 11 ++ 3 files changed, 105 insertions(+), 77 deletions(-) diff --git a/src/gausskernel/storage/replication/syncrep.cpp b/src/gausskernel/storage/replication/syncrep.cpp index a9c6c1ecc..99720b66f 100644 --- a/src/gausskernel/storage/replication/syncrep.cpp +++ b/src/gausskernel/storage/replication/syncrep.cpp @@ -82,12 +82,12 @@ static void SyncRepNotifyComplete(); static int SyncRepGetStandbyPriority(void); #ifndef ENABLE_MULTIPLE_NODES -static bool SyncRepGetCatchupRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr); +static bool SyncRepGetSyncLeftTime(XLogRecPtr XactCommitLSN, TimestampTz* leftTime); #endif -static void SyncRepGetOldestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr, - List* sync_standbys); -static void SyncRepGetNthLatestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr, - List* sync_standbys, uint8 nth); +static void SyncRepGetOldestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, + XLogRecPtr* replayPtr, List* sync_standbys); +static void SyncRepGetNthLatestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, + XLogRecPtr* replayPtr, List* sync_standbys, uint8 nth); #ifdef USE_ASSERT_CHECKING @@ -98,12 +98,6 @@ static List *SyncRepGetSyncStandbysPriority(bool *am_sync, List** catchup_standb static List *SyncRepGetSyncStandbysQuorum(bool *am_sync, List** catchup_standbys = NULL); static int cmp_lsn(const void *a, const void *b); -/* - * Time for synchronous each lsn, which is an empirical value. - * Unit is microsecond. - */ -#define SYN_TIME_PER_XLOG 0.00225 - #define CATCHUP_XLOG_DIFF(ptr1, ptr2, amount) \ XLogRecPtrIsInvalid(ptr1) ? false : (XLByteDifference(ptr2, ptr1) < amount) @@ -113,51 +107,33 @@ static int cmp_lsn(const void *a, const void *b); * Return true if it is need to wait for catching up(synchronous replication), * return false if don't wait for catching up. */ -bool SynRepWaitCatchup(XLogRecPtr XactCommitLSN, int mode) +bool SynRepWaitCatchup(XLogRecPtr XactCommitLSN) { #ifndef ENABLE_MULTIPLE_NODES - /* When most_available_sync is off and num_sync > 1, return. */ + /* + * When most_available_sync is off or catchup2normal_wait_time is not set, + * return true. + */ if (!t_thrd.walsender_cxt.WalSndCtl->most_available_sync || - t_thrd.syncrep_cxt.SyncRepConfig->num_sync > 1 || g_instance.attr.attr_storage.catchup2normal_wait_time < 0) { return true; } - static const uint64 maxXlogDiffAmount = - (uint64)floor(g_instance.attr.attr_storage.catchup2normal_wait_time * 1000 / SYN_TIME_PER_XLOG); - XLogRecPtr receivePtr; - XLogRecPtr writePtr; - XLogRecPtr flushPtr; - XLogRecPtr replayPtr; + static const TimestampTz catchup2normalWaitTime = g_instance.attr.attr_storage.catchup2normal_wait_time * 1000; + TimestampTz syncLeftTime = 0; /* - * if the numberof sync standbys is more than synchronous_standby_names - * specifies or there is no sync standbys in catching up, wait for + * SyncRepGetSyncLeftTime() return false means that there is at lease one + * sync standby, or no standby in catchup, so it is need to wait for * synchronous replication. */ - if (!SyncRepGetCatchupRecPtr(&receivePtr, &writePtr, &flushPtr, &replayPtr)) { + if (!SyncRepGetSyncLeftTime(XactCommitLSN, &syncLeftTime)) { return true; } - - if (maxXlogDiffAmount == 0) { + if (syncLeftTime == 0 || syncLeftTime > catchup2normalWaitTime) { return false; } - - /* if catchup will be finished soon, wait for synchronous replication. */ - switch (mode) { - case SYNC_REP_WAIT_RECEIVE: - return CATCHUP_XLOG_DIFF(receivePtr, XactCommitLSN, maxXlogDiffAmount); - case SYNC_REP_WAIT_WRITE: - return CATCHUP_XLOG_DIFF(writePtr, XactCommitLSN, maxXlogDiffAmount); - case SYNC_REP_WAIT_FLUSH: - return CATCHUP_XLOG_DIFF(flushPtr, XactCommitLSN, maxXlogDiffAmount); - case SYNC_REP_WAIT_APPLY: - return CATCHUP_XLOG_DIFF(replayPtr, XactCommitLSN, maxXlogDiffAmount); - default: - return true; - } #endif - return true; } @@ -205,7 +181,7 @@ void SyncRepWaitForLSN(XLogRecPtr XactCommitLSN) if (!t_thrd.walsender_cxt.WalSndCtl->sync_standbys_defined || XLByteLE(XactCommitLSN, t_thrd.walsender_cxt.WalSndCtl->lsn[mode]) || t_thrd.walsender_cxt.WalSndCtl->sync_master_standalone || - !SynRepWaitCatchup(XactCommitLSN, mode)) { + !SynRepWaitCatchup(XactCommitLSN)) { LWLockRelease(SyncRepLock); RESUME_INTERRUPTS(); return; @@ -681,40 +657,47 @@ bool SyncRepGetSyncRecPtr(XLogRecPtr *receivePtr, XLogRecPtr *writePtr, XLogRecP #ifndef ENABLE_MULTIPLE_NODES /* - * Calculate the synced Receive, Write, Flush and Replay positions among sync standbys - * in catching up. - * - * Return false if the number of sync standbys is more than synchronous_standby_names - * specifies or there is no sync standbys in catching up. Otherwise return true and - * store the positions into *receivePtr, *writePtr, *flushPtr and *replayPtr. + * Obtains the remaining time for synchronizing to the sync standby. + * + * If there is at lease one sync standby, or no standby in catchup, no need + * to consider standby in catchup, return false. Otherwise return true. */ -static bool SyncRepGetCatchupRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr) +static bool SyncRepGetSyncLeftTime(XLogRecPtr XactCommitLSN, TimestampTz* leftTime) { List* sync_standbys = NIL; List* catchup_standbys = NIL; bool am_sync = false; - *receivePtr = InvalidXLogRecPtr; - *writePtr = InvalidXLogRecPtr; - *flushPtr = InvalidXLogRecPtr; - *replayPtr = InvalidXLogRecPtr; + ListCell* cell = NULL; + *leftTime = 0; /* Get standbys that are considered as synchronous at this moment. */ sync_standbys = SyncRepGetSyncStandbys(&am_sync, &catchup_standbys); - /* Quick exit if there are enough synchronous standbys or no standby in catching up. */ - if ((t_thrd.syncrep_cxt.SyncRepConfig != NULL && - t_thrd.syncrep_cxt.SyncRepConfig->num_sync <= list_length(sync_standbys)) || - list_length(catchup_standbys) == 0) { + + /* Skip here if there is at lease one sync standby, or no standby in catchup. */ + if (list_length(sync_standbys) > 0 || list_length(catchup_standbys) == 0) { list_free(sync_standbys); list_free(catchup_standbys); return false; } - if (t_thrd.syncrep_cxt.SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY) { - SyncRepGetOldestSyncRecPtr(receivePtr, writePtr, flushPtr, replayPtr, catchup_standbys); - } else { - SyncRepGetNthLatestSyncRecPtr( - receivePtr, writePtr, flushPtr, replayPtr, catchup_standbys, - t_thrd.syncrep_cxt.SyncRepConfig->num_sync - list_length(sync_standbys)); + + /* + * Scan through all sync standbys and calculate the left time + * for sync. + */ + foreach (cell, catchup_standbys) { + WalSnd* walsnd = &t_thrd.walsender_cxt.WalSndCtl->walsnds[lfirst_int(cell)]; + SpinLockAcquire(&walsnd->mutex); + TimestampTz syncNeededTime; + XLogRecPtr xrp = walsnd->receive; + double rate = walsnd->catchupRate; + SpinLockRelease(&walsnd->mutex); + + syncNeededTime = (TimestampTz)(rate * XLByteDifference(XactCommitLSN, xrp)); + if (*leftTime == 0 || *leftTime > syncNeededTime) { + *leftTime = syncNeededTime; + } } + list_free(sync_standbys); list_free(catchup_standbys); return true; @@ -724,8 +707,8 @@ static bool SyncRepGetCatchupRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr /* * Calculate the oldest Write, Flush and Apply positions among sync standbys. */ -static void SyncRepGetOldestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr, - List* sync_standbys) +static void SyncRepGetOldestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, + XLogRecPtr* replayPtr, List* sync_standbys) { ListCell *cell = NULL; @@ -762,8 +745,8 @@ static void SyncRepGetOldestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* write * Calculate the Nth latest Write, Flush and Apply positions among sync * standbys. */ -static void SyncRepGetNthLatestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, XLogRecPtr* replayPtr, - List* sync_standbys, uint8 nth) +static void SyncRepGetNthLatestSyncRecPtr(XLogRecPtr* receivePtr, XLogRecPtr* writePtr, XLogRecPtr* flushPtr, + XLogRecPtr* replayPtr, List* sync_standbys, uint8 nth) { ListCell *cell = NULL; XLogRecPtr *receive_array = NULL; @@ -1134,9 +1117,9 @@ static List *SyncRepGetSyncStandbysQuorum(bool *am_sync, List** catchup_standbys if (walsnd->sync_standby_priority == 0) continue; - if (walsnd->state == WALSNDSTATE_CATCHUP && catchup_standbys != NULL) { + if ((walsnd->state == WALSNDSTATE_CATCHUP || walsnd->peer_state == CATCHUP_STATE) && + catchup_standbys != NULL) { *catchup_standbys = lappend_int(*catchup_standbys, i); - continue; } /* Must have a valid flush position */ @@ -1206,9 +1189,9 @@ static List *SyncRepGetSyncStandbysPriority(bool *am_sync, List** catchup_standb if (this_priority == 0) continue; - if (walsnd->state == WALSNDSTATE_CATCHUP && catchup_standbys != NULL) { + if ((walsnd->state == WALSNDSTATE_CATCHUP || walsnd->peer_state == CATCHUP_STATE) && + catchup_standbys != NULL) { *catchup_standbys = lappend_int(*catchup_standbys, i); - continue; } /* Must have a valid flush position */ diff --git a/src/gausskernel/storage/replication/walsender.cpp b/src/gausskernel/storage/replication/walsender.cpp index 322d005e2..13bc3992f 100644 --- a/src/gausskernel/storage/replication/walsender.cpp +++ b/src/gausskernel/storage/replication/walsender.cpp @@ -117,6 +117,11 @@ long g_logical_slot_sleep_time = 0; #define AmWalSenderToDummyStandby() (t_thrd.walsender_cxt.MyWalSnd->sendRole == SNDROLE_PRIMARY_DUMMYSTANDBY) #define AmWalSenderOnDummyStandby() (t_thrd.walsender_cxt.MyWalSnd->sendRole == SNDROLE_DUMMYSTANDBY_STANDBY) +/* + * calculate catchup late every 1000ms + */ +#define CALCULATE_CATCHUP_RATE_TIME 1000 + /* Statistics for log control */ static const int MICROSECONDS_PER_SECONDS = 1000000; static const int MILLISECONDS_PER_SECONDS = 1000; @@ -208,6 +213,7 @@ static void WalSndRefreshPercentCountStartLsn(XLogRecPtr currentMaxLsn, XLogRecP static void set_xlog_location(ServerMode local_role, XLogRecPtr* sndWrite, XLogRecPtr* sndFlush, XLogRecPtr* sndReplay); static void ProcessArchiveFeedbackMessage(void); static void WalSndArchiveXlog(XLogRecPtr targetLsn, int sub_term); +static void CalCatchupRate(); char *DataDir = "."; @@ -3325,13 +3331,8 @@ static int WalSndLoop(WalSndSendDataCallback send_data) } } } else { - if (t_thrd.walsender_cxt.MyWalSnd->state == WALSNDSTATE_STREAMING && - !XLByteLT(t_thrd.walsender_cxt.catchup_threshold, - INT2UINT64(g_instance.attr.attr_storage.MaxSendSize) * 1024)) { - ereport(DEBUG1, (errmsg("standby \"%s\" has now caught up with primary", - u_sess->attr.attr_common.application_name))); - WalSndSetState(WALSNDSTATE_CATCHUP); - t_thrd.walsender_cxt.catchup_threshold = 0; + if (t_thrd.walsender_cxt.MyWalSnd->state == WALSNDSTATE_CATCHUP) { + CalCatchupRate(); } } @@ -3618,6 +3619,9 @@ static void InitWalSnd(void) walsnd->log_ctrl.pre_rate2 = 0; walsnd->log_ctrl.prev_RPO = -1; walsnd->log_ctrl.current_RPO = -1; + walsnd->lastCalTime = 0; + walsnd->lastCalWrite = InvalidXLogRecPtr; + walsnd->catchupRate = 0; walsnd->log_ctrl.just_keep_alive = false; SpinLockRelease(&walsnd->mutex); /* don't need the lock anymore */ @@ -5468,3 +5472,33 @@ XLogSegNo WalGetSyncCountWindow(void) { return (XLogSegNo)(uint32)u_sess->attr.attr_storage.wal_keep_segments; } + +/* + * Calculate catchup rate of standby to estimate how long + * the standby will be caught up with primary. + */ +static void CalCatchupRate() { + if (g_instance.attr.attr_storage.catchup2normal_wait_time < 0) { + return; + } + volatile WalSnd *walsnd = t_thrd.walsender_cxt.MyWalSnd; + + SpinLockAcquire(&walsnd->mutex); + XLogRecPtr write = walsnd->write; + TimestampTz now = GetCurrentTimestamp(); + + if (XLByteEQ(walsnd->lastCalWrite, InvalidXLogRecPtr)) { + walsnd->lastCalWrite = write; + walsnd->lastCalTime = now; + SpinLockRelease(&walsnd->mutex); + return; + } + if (TimestampDifferenceExceeds(walsnd->lastCalTime, now, CALCULATE_CATCHUP_RATE_TIME)) { + double tempRate = (double)(now - walsnd->lastCalTime) / + (double)XLByteDifference(write, walsnd->lastCalWrite); + walsnd->catchupRate = walsnd->catchupRate == 0 ? tempRate : (walsnd->catchupRate + tempRate) / 2; + walsnd->lastCalWrite = write; + walsnd->lastCalTime = now; + } + SpinLockRelease(&walsnd->mutex); +} diff --git a/src/include/replication/walsender_private.h b/src/include/replication/walsender_private.h index feba21b35..a097fb212 100644 --- a/src/include/replication/walsender_private.h +++ b/src/include/replication/walsender_private.h @@ -120,6 +120,17 @@ typedef struct WalSnd { LogCtrlData log_ctrl; unsigned int archive_flag; Latch* arch_latch; + + /* + * lastCalTime is last time calculating catchupRate, and lastCalWrite + * is last calculating write lsn. + */ + TimestampTz lastCalTime; + XLogRecPtr lastCalWrite; + /* + * Time needed for synchronous per xlog while catching up. + */ + double catchupRate; } WalSnd; extern THR_LOCAL WalSnd* MyWalSnd;