极致RTO优化

This commit is contained in:
zhaowenhao 2021-05-29 14:45:28 +08:00
parent 7d89cf9b71
commit 29cda19e0a
3 changed files with 39 additions and 30 deletions

View File

@ -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)
{

View File

@ -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",

View File

@ -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;