SysCache总览
| 缩略语 | 全称 | 中文名 | 备注 | 调用接口 |
|---|---|---|---|---|
| CatCache | Catalog Cache | 系统表数据缓存 | 缓存系统表行 | SearchSysCache |
| RelCache | Relation Cache | 表结构缓存 | 表定义,表Oid到RelationData的映射,DDL语句的并发会刷新绝大部分字段 | RelationIdGetRelation |
| PartitionCache | Partition Cache | 分区表缓存 | 分区Oid到PartitionData的映射 | PartitionIdGetPartition |
| RelMap | Relation filenode map | 表filenode映射 | 表Oid到filenode的映射 | RelationMapOidToFilenode RelationMapFilenodeToOid |
缓存失效
缓存根据存放的位置可分为local和global,顾名思义是本地和全局,在openGauss上就是线程级和进程级缓存的区别。以存放RelationData的RelCache为例,如果一个线程执行了DDL语句,导致表对应的RelationData发生了变更,例如字段信息、分区信息等,这意味着其他线程以及进程上已有的RelationData已经过时,需要更新或删除。那么通知其何时刷新则是通过失效消息机制来实现。
失效消息机制的主体是SI(SharedInvalidation) Message Queue,它是一个全局唯一的循环队列,存放失效消息。同样全局维护一个ProcState数组,每启动一个后台线程,都会在上面注册一个ProcState,记录当前队列的读取位置。
失效消息的全路径可分为注册、发送以及处理。注册是在事务还没提交前,先在本地保存失效消息,提交之后才发送到全局消息队列。以drop partition为例,分别走读这三个阶段的代码。
CREATE TABLE IF NOT EXISTS t1 ( id varchar(30), now_date varchar(300), create_time TIMESTAMP NOT NULL ) WITH(ORIENTATION = ROW, COMPRESSION = NO) PARTITION BY RANGE (now_date) ( PARTITION part_20240327 VALUES LESS THAN ('20240327'), PARTITION part_20240328 VALUES LESS THAN ('20240328'), PARTITION part_20240329 VALUES LESS THAN ('20240329'), PARTITION part_20240330 VALUES LESS THAN ('20240330')); -- 会话一,查询该表,加载RelationData到本地缓存 select * from t1 where id = '1'; -- 会话二,drop partition alter table t1 drop partition part_20240329;执行DDL的线程:
其他线程:
注册失效消息
删分区的操作有两处注册,一个是fastDropPartition,一个是UpdatePgObjectMtime。
void fastDropPartition(Relation rel, ..., bool sendInvalid) { ... /* 发送失效消息 */ if (sendInvalid) { if (RelationIsPartitionOfSubPartitionTable(rel)) { /* 二级分区 */ CacheInvalidatePartcacheByPartid(rel->rd_id); } else { CacheInvalidateRelcache(rel); } } /* command计数器加一,同时处理失效消息,失效本地缓存 */ CommandCounterIncrement(); }void CacheInvalidateRelcache(Relation relation) { Oid databaseId; Oid relationId; relationId = RelationGetRelid(relation); if (relation->rd_rel->relisshared) { /* 共享表,不属于某个database,如pg_database */ databaseId = InvalidOid; } else { databaseId = u_sess->proc_cxt.MyDatabaseId; } RegisterRelcacheInvalidation(databaseId, relationId); }static void RegisterRelcacheInvalidation(Oid dbId, Oid relId) { /* 添加失效消息到线程级别的列表中 */ AddRelcacheInvalidationMessage(&GetInvalCxt()->transInvalInfo->CurrentCmdInvalidMsgs, dbId, relId); (void)GetCurrentCommandId(true); /* relcache存在一个initfile,启动数据库时会从该文件中读取写入到relcache,缓存系统表。用户表不在里面。如果是系统表的ddl,需要将该文件置为无效 */ if (RelationIdIsInInitFile(relId)) { u_sess->inval_cxt.transInvalInfo->RelcacheInitFileInval = true; } }CurrentCmdInvalidMsgs保存线程当前cmd注册的失效消息,因为还没提交,就暂时不发到全局的消息列表。通过CommandCounterIncrement,会将CurrentCmdInvalidMsgs的消息放到PriorCmdInvalidMsgs中,并将前者置为NULL。
static void AddRelcacheInvalidationMessage(InvalidationListHeader* hdr, Oid dbId, Oid relId) { SharedInvalidationMessage msg; /* 遍历当前的列表,如果有重复的,直接return */ ProcessMessageList(hdr->rclist, if (msg->rc.id == SHAREDINVALRELCACHE_ID && (msg->rc.relId == relId || msg->rc.relId == InvalidOid)) return); /* OK, add the item */ msg.rc.id = SHAREDINVALRELCACHE_ID; msg.rc.dbId = dbId; msg.rc.relId = relId; AddInvalidationMessage(&hdr->rclist, &msg); }static void AddInvalidationMessage(InvalidationChunk** listHdr, SharedInvalidationMessage* msg) { InvalidationChunk* chunk = *listHdr; if (chunk == NULL) { /* 首次进来,创建首块。FIRSTCHUNKSIZE-1是因为InvalidationChunk自带一个 */ #define FIRSTCHUNKSIZE 32 chunk = (InvalidationChunk*)MemoryContextAlloc(t_thrd.mem_cxt.cur_transaction_mem_cxt, sizeof(InvalidationChunk) + (FIRSTCHUNKSIZE - 1) * sizeof(SharedInvalidationMessage)); chunk->nitems = 0; chunk->maxitems = FIRSTCHUNKSIZE; chunk->next = *listHdr; *listHdr = chunk; } else if (chunk->nitems >= chunk->maxitems) { /* 如果超过maxitems,再扩展一个chunk,组成链表 */ int chunksize = 2 * chunk->maxitems; chunk = (InvalidationChunk*)MemoryContextAlloc(t_thrd.mem_cxt.cur_transaction_mem_cxt, sizeof(InvalidationChunk) + (chunksize - 1) * sizeof(SharedInvalidationMessage)); chunk->nitems = 0; chunk->maxitems = chunksize; chunk->next = *listHdr; *listHdr = chunk; } /* 添加消息到当前chunk */ chunk->msgs[chunk->nitems] = *msg; chunk->nitems++; }注册完成后,fastDropPartition最后还调用了CommandCounterIncrement,该函数处理自己产生的失效消息,失效本地缓存,保证后续本线程的一致性。
void CommandEndInvalidationMessages(void) { knl_u_inval_context *inval_cxt = GetInvalCxt(); if (inval_cxt->transInvalInfo == NULL) { return; } /* 处理全局缓存 */ ProcessInvalidationMessagesMulti( &inval_cxt->transInvalInfo->CurrentCmdInvalidMsgs, GlobalExecuteSharedInvalidMessages); /* 处理本地缓存 */ ProcessInvalidationMessages( &inval_cxt->transInvalInfo->CurrentCmdInvalidMsgs, LocalExecuteThreadAndSessionInvalidationMessage); /* 将CurrentCmdInvalidMsgs列表加到PriorCmdInvalidMsgs */ AppendInvalidationMessages(&inval_cxt->transInvalInfo->PriorCmdInvalidMsgs, &inval_cxt->transInvalInfo->CurrentCmdInvalidMsgs); }void GlobalExecuteSharedInvalidMessages(const SharedInvalidationMessage* msgs, int n) { /* 没开gsc不需要处理 */ if (!EnableLocalSysCache()) { return; } /* 非THREADPOOL_STREAM和bgworker,传入is_commit为false */ GlobalInvalidSharedInvalidMessages(msgs, n, IS_THREAD_POOL_STREAM || IsBgWorkerProcess()); }static void GlobalInvalidSharedInvalidMessages(const SharedInvalidationMessage* msgs, int n, bool is_commit) { Assert(EnableGlobalSysCache()); for (int i = 0; i < n; i++) { SharedInvalidationMessage *msg = (SharedInvalidationMessage *)(msgs + i); if (msg->id >= 0) { t_thrd.lsc_cxt.lsc->systabcache.CacheIdHashValueInvalidateGlobal(msg->cc.dbId, msg->cc.id, msg->cc.hashValue, is_commit); } else if (msg->id == SHAREDINVALCATALOG_ID) { t_thrd.lsc_cxt.lsc->systabcache.CatalogCacheFlushCatalogGlobal(msg->cat.dbId, msg->cat.catId, is_commit); } else if (msg->id == SHAREDINVALRELCACHE_ID) { /* 处理relcache相关的失效消息 */ t_thrd.lsc_cxt.lsc->tabdefcache.InvalidateGlobalRelation(msg->rc.dbId, msg->rc.relId, is_commit); } else if (msg->id == SHAREDINVALPARTCACHE_ID) { t_thrd.lsc_cxt.lsc->partdefcache.InvalidateGlobalPartition(msg->pc.dbId, msg->pc.partId, is_commit); } /* global relmap are backups of relmapfile, so no need to deal with the msg */ } }void LocalTabDefCache::InvalidateGlobalRelation(Oid db_id, Oid rel_oid, bool is_commit) { if (unlikely(db_id == InvalidOid && rel_oid == InvalidOid)) { return; } if (!is_commit) { /* 非提交场景下,仅将rel_oid存入invalid_entries数组。后续本线程从global cache获取时,如果invalid_entries包含该oid,直接return NULL */ invalid_entries.InsertInvalidDefValue(rel_oid); return; } if (db_id == InvalidOid) { t_thrd.lsc_cxt.lsc->GetSharedTabDefCache()->Invalidate(db_id, rel_oid); } else if (m_global_tabdefcache == NULL) { Assert(!m_is_inited_phase3); Assert(CheckMyDatabaseMatch()); GlobalSysDBCacheEntry *entry = g_instance.global_sysdbcache.FindTempGSCEntry(db_id); if (entry == NULL) { return; } entry->m_tabdefCache->Invalidate(db_id, rel_oid); g_instance.global_sysdbcache.ReleaseTempGSCEntry(entry); } else { Assert(CheckMyDatabaseMatch()); Assert(m_db_id == t_thrd.lsc_cxt.lsc->my_database_id); Assert(m_db_id == db_id); /* 提交场景下失效全局的relation */ m_global_tabdefcache->Invalidate(db_id, rel_oid); } }发送失效消息
提交时才将失效消息存入全局的列表。
void AtEOXact_Inval(bool isCommit) { knl_u_inval_context *inval_cxt = GetInvalCxt(); if (isCommit) { Assert(inval_cxt->transInvalInfo != NULL && inval_cxt->transInvalInfo->parent == NULL); /* 获取RelcacheInit锁,删除initfile */ if (inval_cxt->transInvalInfo->RelcacheInitFileInval) { RelationCacheInitFilePreInvalidate(); } /* 再将最近一次CurrentCmdInvalidMsgs加入PriorCmdInvalidMsgs */ AppendInvalidationMessages(&inval_cxt->transInvalInfo->PriorCmdInvalidMsgs, &inval_cxt->transInvalInfo->CurrentCmdInvalidMsgs); /* 发送PriorCmdInvalidMsgs里的失效消息 */ ProcessInvalidationMessagesMulti( &inval_cxt->transInvalInfo->PriorCmdInvalidMsgs, SendSharedInvalidMessages); /* 释放RelcacheInit锁 */ if (inval_cxt->transInvalInfo->RelcacheInitFileInval) { RelationCacheInitFilePostInvalidate(); } } else if (inval_cxt->transInvalInfo != NULL) { /* Must be at top of stack */ Assert(inval_cxt->transInvalInfo->parent == NULL); ProcessInvalidationMessages( &inval_cxt->transInvalInfo->PriorCmdInvalidMsgs, LocalExecuteThreadAndSessionInvalidationMessage); } /* Need not free anything explicitly */ inval_cxt->transInvalInfo = NULL; }static void ProcessInvalidationMessagesMulti( InvalidationListHeader* hdr, void (*func)(const SharedInvalidationMessage* msgs, int n)) { ProcessMessageListMulti(hdr->cclist, func(msgs, n)); /* 发送relcache相关,每次处理一个chunk */ ProcessMessageListMulti(hdr->rclist, func(msgs, n)); ProcessMessageListMulti(hdr->pclist, func(msgs, n)); }全局缓存由执行ddl的线程失效,其他线程的local缓存由自身处理失效消息来失效。
void SendSharedInvalidMessages(const SharedInvalidationMessage* msgs, int n) { if (EnableGlobalSysCache()) { GlobalInvalidSharedInvalidMessages(msgs, n, true); } /* 没开gsc也需要失效其他线程的缓存 */ SIInsertDataEntries(msgs, n); if (ENABLE_GPC && g_instance.plan_cache != NULL) { g_instance.plan_cache->InvalMsg(msgs, n); } }void SIInsertDataEntries(const SharedInvalidationMessage* data, int n) { /* 全局失效buffer放在共享内存中,以循环缓冲区的形式 */ SISeg* segP = t_thrd.shemem_ptr_cxt.shmInvalBuffer; /* 每次仅处理WRITE_QUANTUM条消息,避免持锁时间过长 */ while (n > 0) { int nthistime = Min(n, WRITE_QUANTUM); int numMsgs; int max; int i; n -= nthistime; LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE); /* 如果当前缓冲队列中的条目数过多,需要做清理 */ for (;;) { numMsgs = segP->maxMsgNum - segP->minMsgNum; if (numMsgs + nthistime > MAXNUMMESSAGES || numMsgs >= segP->nextThreshold) { SICleanupQueue(true, nthistime); } else { break; } } /* 插入循环缓冲区中 */ max = segP->maxMsgNum; while (nthistime-- > 0) { segP->buffer[max % MAXNUMMESSAGES] = *data++; max++; } /* 更新maxMsgNum,加自旋锁 */ { /* use volatile pointer to prevent code rearrangement */ volatile SISeg* vsegP = segP; SpinLockAcquire(&vsegP->msgnumLock); vsegP->maxMsgNum = max; SpinLockRelease(&vsegP->msgnumLock); } /* 遍历procState,将hasMessages置为true,表示有消息未处理 */ for (i = 0; i < segP->lastBackend; i++) { ProcState* stateP = &segP->procState[i]; if (stateP->procPid != 0) { stateP->hasMessages = true; } } LWLockRelease(SInvalWriteLock); } }void SICleanupQueue(bool callerHasWriteLock, int minFree) { SISeg* segP = t_thrd.shemem_ptr_cxt.shmInvalBuffer; int min, minsig, lowbound, numMsgs, i; ProcState* needSig = NULL; /* Lock out all writers and readers */ if (!callerHasWriteLock) { LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE); } LWLockAcquire(SInvalReadLock, LW_EXCLUSIVE); /* 重新计算minMsgNum,即所有后台线程中最小nextMsgNum */ min = segP->maxMsgNum; /* 落后于minsig才发送catchup信号 */ minsig = min - SIG_THRESHOLD; /* 保证清理完后有足够空间的下界 */ lowbound = min - MAXNUMMESSAGES + minFree; for (i = 0; i < segP->lastBackend; i++) { ProcState* stateP = &segP->procState[i]; int n = stateP->nextMsgNum; /* 忽略不活跃的、resetState已经被置为true的 */ if (stateP->procPid == 0 || stateP->resetState || stateP->sendOnly) { continue; } /* 如果小于下界,直接设置resetState为true */ if (n < lowbound) { stateP->resetState = true; /* no point in signaling him ... */ continue; } /* 计算最小nextMsgNum */ if (n < min) { min = n; } /* 找一个不超过下界范围内最落后的,发送信号 */ if ((i < segP->maxreserveBackends) && (n < minsig) && (!stateP->signaled)) { minsig = n; needSig = stateP; } } segP->minMsgNum = min; /* 当minMsgNum过大,需要减掉一个固定值,防止发生整型回绕 */ if (min >= MSGNUMWRAPAROUND) { segP->minMsgNum -= MSGNUMWRAPAROUND; segP->maxMsgNum -= MSGNUMWRAPAROUND; for (i = 0; i < segP->lastBackend; i++) { /* we don't bother skipping inactive entries here */ if (segP->procState[i].procPid != 0) segP->procState[i].nextMsgNum -= MSGNUMWRAPAROUND; } } /* 计算有多少msg,并且决定后面再次调用SICleanupQueue的阈值 */ numMsgs = segP->maxMsgNum - segP->minMsgNum; if (numMsgs < CLEANUP_MIN) { segP->nextThreshold = CLEANUP_MIN; } else { segP->nextThreshold = (numMsgs / CLEANUP_QUANTUM + 1) * CLEANUP_QUANTUM; } /* 最后给needSig发送catchup信号 */ if (needSig != NULL) { ThreadId his_pid = needSig->procPid; BackendId his_backendId = (needSig - &segP->procState[0]) + 1; needSig->signaled = true; LWLockRelease(SInvalReadLock); LWLockRelease(SInvalWriteLock); ereport(DEBUG4, (errmsg("sending sinval catchup signal to ThreadId %lu", his_pid))); SendProcSignal(his_pid, PROCSIG_CATCHUP_INTERRUPT, his_backendId); if (callerHasWriteLock) { LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE); } } else { LWLockRelease(SInvalReadLock); if (!callerHasWriteLock) { LWLockRelease(SInvalWriteLock); } } }处理失效消息
主函数是AcceptInvalidationMessages,调用点比较多,全局搜大概有40来处,个人觉得重点的是开启事务(AtStart_Cache)和加锁(LockRelation\LockRelationOid…)
void AcceptInvalidationMessages() { AssertEreport(t_thrd.inval_msg_cxt.b_can_not_process == false, MOD_OPT, "unable to accept invalidation messages"); /* 处理失效消息也存在递归调用,例如下图 */ if (!DeepthInAcceptInvalidationMessageNotZero()) { t_thrd.rc_cxt.rcNum = 0; } if (EnableLocalSysCache()) { u_sess->pcache_cxt.gpc_remote_msg = true; knl_u_inval_context *inval_cxt = &t_thrd.lsc_cxt.lsc->inval_cxt; /* 递归层数+1 */ ++inval_cxt->DeepthInAcceptInvalidationMessage; if (!IS_THREAD_POOL_WORKER) { ReceiveSharedInvalidMessages(LocalExecuteThreadAndSessionInvalidationMessage, InvalidateSystemCaches, false); } else { /* 线程池模式 */ ReceiveSharedInvalidMessages(LocalExecuteThreadInvalidationMessage, InvalidateThreadSystemCaches, false); u_sess->pcache_cxt.gpc_remote_msg = false; TestCodeToForceCacheFlushes(); --inval_cxt->DeepthInAcceptInvalidationMessage; u_sess->pcache_cxt.gpc_remote_msg = true; inval_cxt = &u_sess->inval_cxt; ++inval_cxt->DeepthInAcceptInvalidationMessage; ReceiveSharedInvalidMessages(LocalExecuteSessionInvalidationMessage, InvalidateSessionSystemCaches, true); } u_sess->pcache_cxt.gpc_remote_msg = false; TestCodeToForceCacheFlushes(); --inval_cxt->DeepthInAcceptInvalidationMessage; return; } /* 没开gsc的case */ u_sess->pcache_cxt.gpc_remote_msg = true; ++u_sess->inval_cxt.DeepthInAcceptInvalidationMessage; ReceiveSharedInvalidMessages(LocalExecuteInvalidationMessage, InvalidateSystemCaches, false); u_sess->pcache_cxt.gpc_remote_msg = false; TestCodeToForceCacheFlushes(); --u_sess->inval_cxt.DeepthInAcceptInvalidationMessage; }在重建RelationData的过程中,需要扫描系统表,加锁之后再次调用AcceptInvalidationMessages。
void ReceiveSharedInvalidMessages(void (*invalFunction)(SharedInvalidationMessage* msg), void (*resetFunction)(void), bool worksession) { knl_u_inval_context *inval_cxt; if (!EnableLocalSysCache()) { inval_cxt = &u_sess->inval_cxt; } else if (worksession) { inval_cxt = &u_sess->inval_cxt; } else { inval_cxt = &t_thrd.lsc_cxt.lsc->inval_cxt; } /* 处理外层调用获取的消息 */ while (inval_cxt->nextmsg < inval_cxt->nummsgs) { SharedInvalidationMessage msg = inval_cxt->messages[inval_cxt->nextmsg++]; if (SkipRedundantInvalMsg(&msg)) { continue; } inval_cxt->SIMCounter++; invalFunction(&msg); } do { int getResult; inval_cxt->nextmsg = inval_cxt->nummsgs = 0; /* 获取更多未处理的消息,返回数目 */ getResult = SIGetDataEntries(inval_cxt->messages, MAXINVALMSGS, worksession); if (getResult < 0) { /* got a reset message */ ereport(DEBUG4, (errmsg("cache state reset"))); inval_cxt->SIMCounter++; /* 返回-1,表示将所有缓存失效,调用reset函数 */ resetFunction(); break; /* nothing more to do */ } /* 重置nextmsg为0 */ inval_cxt->nextmsg = 0; inval_cxt->nummsgs = getResult; while (inval_cxt->nextmsg < inval_cxt->nummsgs) { SharedInvalidationMessage msg = inval_cxt->messages[inval_cxt->nextmsg++]; /* 记录处理的失效消息数 */ inval_cxt->SIMCounter++; invalFunction(&msg); } /* 如果SIGetDataEntries获取的消息数填满了messages,则说明可能还有消息没处理完,接着处理 */ } while (inval_cxt->nummsgs == MAXINVALMSGS); /* 当前已经追上了,如果catchupInterruptPending为true,置为false,并且调用SICleanupQueue,目的是给下一个最落后的线程发送catchup信号 */ if (catchupInterruptPending) { catchupInterruptPending = false; ereport(DEBUG4, (errmsg("sinval catchup complete, cleaning queue"))); SICleanupQueue(false, 0); } }int SIGetDataEntries(SharedInvalidationMessage* data, int datasize, bool worksession) { SISeg* segP = NULL; ProcState* stateP = NULL; int max; int n; segP = t_thrd.shemem_ptr_cxt.shmInvalBuffer; if (IS_THREAD_POOL_WORKER) { /* 线程池模式 */ if (EnableLocalSysCache() && !worksession) { /* GSC is on, and this is a thread pool worker, fetch inval msg slot by backendid */ stateP = &segP->procState[t_thrd.proc_cxt.MyBackendId - 1]; } else { /* this is a thread pool worker, fetch inval msg slot by session index */ stateP = &segP->procState[u_sess->session_ctr_index]; } } else { /* 根据backendId获取ProcState */ stateP = &segP->procState[t_thrd.proc_cxt.MyBackendId - 1]; } /* 快速返回 */ if (!stateP->hasMessages) { return 0; } LWLockAcquire(SInvalReadLock, LW_SHARED); /* 将hasMessages置为false */ stateP->hasMessages = false; /* 获取自旋锁,得到当前的maxMsgNum */ { /* use volatile pointer to prevent code rearrangement */ volatile SISeg* vsegP = segP; SpinLockAcquire(&vsegP->msgnumLock); max = vsegP->maxMsgNum; SpinLockRelease(&vsegP->msgnumLock); } if (stateP->resetState) { /* 强制reset,表示所有的缓存均失效 */ stateP->nextMsgNum = max; stateP->resetState = false; stateP->signaled = false; LWLockRelease(SInvalReadLock); return -1; } /* 将信息拷贝到data数组中 */ n = 0; while (n < datasize && stateP->nextMsgNum < max) { data[n++] = segP->buffer[stateP->nextMsgNum % MAXNUMMESSAGES]; stateP->nextMsgNum++; } /* 如果读完了所有消息,将signaled置为false;如果还有,将hasMessages置为true */ if (stateP->nextMsgNum >= max) { stateP->signaled = false; } else { stateP->hasMessages = true; } LWLockRelease(SInvalReadLock); return n; }