openGauss_syscache缓存失效机制

📅 发布时间:2026/8/15 8:36:11
openGauss_syscache缓存失效机制
SysCache总览缩略语全称中文名备注调用接口CatCacheCatalog Cache系统表数据缓存缓存系统表行SearchSysCacheRelCacheRelation Cache表结构缓存表定义表Oid到RelationData的映射DDL语句的并发会刷新绝大部分字段RelationIdGetRelationPartitionCachePartition Cache分区表缓存分区Oid到PartitionData的映射PartitionIdGetPartitionRelMapRelation 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 dont 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; }