From 29cda19e0a7dbd5c28491019937ef35f79c852ca Mon Sep 17 00:00:00 2001 From: zhaowenhao <545612025@qq.com> Date: Sat, 29 May 2021 14:45:28 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9E=81=E8=87=B4RTO=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../access/transam/extreme_rto/dispatcher.cpp | 46 ++++++++++--------- .../access/transam/extreme_rto/page_redo.cpp | 10 ++-- .../storage/access/transam/xlog.cpp | 13 ++++-- 3 files changed, 39 insertions(+), 30 deletions(-) diff --git a/src/gausskernel/storage/access/transam/extreme_rto/dispatcher.cpp b/src/gausskernel/storage/access/transam/extreme_rto/dispatcher.cpp index 6912b8272..462d6894b 100644 --- a/src/gausskernel/storage/access/transam/extreme_rto/dispatcher.cpp +++ b/src/gausskernel/storage/access/transam/extreme_rto/dispatcher.cpp @@ -325,20 +325,39 @@ void AllocRecordReadBuffer(XLogReaderState *xlogreader, uint32 privateLen) #endif } +void SendSingalToPageWorker(int signal) +{ + for (uint32 i = 0; i < g_instance.comm_cxt.predo_cxt.totalNum; ++i) { + uint32 state = pg_atomic_read_u32(&(g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadState)); + if (state == PAGE_REDO_WORKER_READY) { + int err = gs_signal_send(g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadId, signal); + if (0 != err) { + ereport(WARNING, (errmsg("Dispatch kill(pid %lu, signal %d) failed: \"%s\",", + g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadId, signal, + gs_strerror(err)))); + } + } + } +} + void HandleStartupInterruptsForExtremeRto() { Assert(AmStartupProcess()); - uint32 triggeredstate = pg_atomic_read_u32(&(g_startupTriggerState)); + uint32 newtriggered = (uint32)CheckForSatartupStatus(); - if (newtriggered != extreme_rto::TRIGGER_NORMAL && triggeredstate != newtriggered) { - ereport(LOG, (errmodule(MOD_REDO), errcode(ERRCODE_LOG), - errmsg("HandleStartupInterruptsForExtremeRto:g_startupTriggerState set from %u to %u", - triggeredstate, newtriggered))); - pg_atomic_write_u32(&(g_startupTriggerState), newtriggered); + if (newtriggered != extreme_rto::TRIGGER_NORMAL) { + uint32 triggeredstate = pg_atomic_read_u32(&(g_startupTriggerState)); + if (triggeredstate != newtriggered) { + ereport(LOG, (errmodule(MOD_REDO), errcode(ERRCODE_LOG), + errmsg("HandleStartupInterruptsForExtremeRto:g_startupTriggerState set from %u to %u", + triggeredstate, newtriggered))); + pg_atomic_write_u32(&(g_startupTriggerState), newtriggered); + } } if (t_thrd.startup_cxt.got_SIGHUP) { t_thrd.startup_cxt.got_SIGHUP = false; + SendSingalToPageWorker(SIGHUP); ProcessConfigFile(PGC_SIGHUP); } @@ -530,21 +549,6 @@ void SetPageWorkStateByThreadId(uint32 threadState) } } -void SendSingalToPageWorker(int signal) -{ - for (uint32 i = 0; i < g_instance.comm_cxt.predo_cxt.totalNum; ++i) { - uint32 state = pg_atomic_read_u32(&(g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadState)); - if (state == PAGE_REDO_WORKER_READY) { - int err = gs_signal_send(g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadId, signal); - if (0 != err) { - ereport(WARNING, (errmsg("Dispatch kill(pid %lu, signal %d) failed: \"%s\",", - g_instance.comm_cxt.predo_cxt.pageRedoThreadStatusList[i].threadId, signal, - gs_strerror(err)))); - } - } - } -} - /* Run from the dispatcher thread. */ static void StopRecoveryWorkers(int code, Datum arg) { diff --git a/src/gausskernel/storage/access/transam/extreme_rto/page_redo.cpp b/src/gausskernel/storage/access/transam/extreme_rto/page_redo.cpp index 61d9dbef7..052e94766 100644 --- a/src/gausskernel/storage/access/transam/extreme_rto/page_redo.cpp +++ b/src/gausskernel/storage/access/transam/extreme_rto/page_redo.cpp @@ -1258,8 +1258,9 @@ void StartupSendFowarder(RedoItem *item) void SendLsnFowarder() { - GetCompletedReadEndPtr(g_redoWorker, &g_GlobalLsnForwarder.record.ReadRecPtr, - &g_GlobalLsnForwarder.record.EndRecPtr); + // update and read in the same thread, so no need atomic operation + g_GlobalLsnForwarder.record.ReadRecPtr = g_redoWorker->lastReplayedReadRecPtr; + g_GlobalLsnForwarder.record.EndRecPtr = g_redoWorker->lastReplayedEndRecPtr; g_GlobalLsnForwarder.record.refcount = get_real_recovery_parallelism() - XLOG_READER_NUM; g_GlobalLsnForwarder.record.isDecode = true; PutRecordToReadQueue(&g_GlobalLsnForwarder.record); @@ -1625,10 +1626,11 @@ void XLogReadPageWorkerMain() PutRecordToReadQueue(xlogreader); xlogreader = newxlogreader; - SetCompletedReadEndPtr(g_redoWorker, xlogreader->ReadRecPtr, xlogreader->EndRecPtr); + g_redoWorker->lastReplayedReadRecPtr = xlogreader->ReadRecPtr; + g_redoWorker->lastReplayedEndRecPtr = xlogreader->EndRecPtr; RedoInterruptCallBack(); - if (CheckForForceFinishRedoTrigger(&term_file)) { + if (force_finish_enabled() && CheckForForceFinishRedoTrigger(&term_file)) { ereport(WARNING, (errmsg("[ForceFinish] force finish triggered in XLogReadPageWorkerMain, ReadRecPtr:%08X/%08X, " "EndRecPtr:%08X/%08X, StandbyMode:%u, startup_processing:%u, dummyStandbyMode:%u", diff --git a/src/gausskernel/storage/access/transam/xlog.cpp b/src/gausskernel/storage/access/transam/xlog.cpp index 0bfe919f4..b3ea96acd 100644 --- a/src/gausskernel/storage/access/transam/xlog.cpp +++ b/src/gausskernel/storage/access/transam/xlog.cpp @@ -16631,17 +16631,20 @@ bool CheckFinishRedoSignal(void) extreme_rto::Enum_TriggeredState CheckForSatartupStatus(void) { - if (CheckForPrimaryTrigger()) { - /* update flag */ + if (t_thrd.startup_cxt.primary_triggered) { + ereport(LOG, (errmsg("received primary request"))); + ResetPrimaryTriggered(); return extreme_rto::TRIGGER_PRIMARY; } - if (CheckForStandbyTrigger()) { + if (t_thrd.startup_cxt.standby_triggered) { + ereport(LOG, (errmsg("received standby request"))); + ResetStandbyTriggered(); return extreme_rto::TRIGGER_STADNBY; } - if (IsFailoverTriggered()) { + if (t_thrd.startup_cxt.failover_triggered) { return extreme_rto::TRIGGER_FAILOVER; } - if (IsSwitchoverTriggered()) { + if (t_thrd.startup_cxt.switchover_triggered) { return extreme_rto::TRIGGER_FAILOVER; } return extreme_rto::TRIGGER_NORMAL;