PostgreSQL 3-3 Buffer

1 Buffer基本介绍 1.1 Buffer解决的问题 实际应用场景中,PostgreSQL可能需要1分钟处理数十万次事务,频繁读写Page。Buffer提供以下2个基本功能: 加载Page至内存:PostgreSQL需要读写数时,都需先从磁盘中,将指定的Page加载至内存中,再在内存中操作Page 缓存常用Page:Page被加载至内存后,还会在内存停留一段时间,降低PostgreSQL直接从磁盘读写Page的频率 1.2 Buffer的基本功能 Buffer缓存Page时,需考虑以下几个场景: 插入Page:大部分事务操作表的Page时,先从Buffer中查找Page,如果Page不在Buffer中,再从磁盘读取Page,将其放入Buffer中,再操作Buffer中的Page 查找Page:访问Page时,通过文件名和Page编号等信息可确定唯一Page 访问Page:所有事务均可访问Buffer,需考虑多事务并发读写同一Pape的场景 淘汰Page:Buffer能缓存的Page数量有限,在Page数量过多时,需淘汰一些Page 落盘Page:Buffer中通常缓存表数据文件的Page,前文提到,事务持久化的标志是WAL日志落盘,因此,Buffer中的Page无需实时落盘,仅需异步落盘即可 2 Buffer的基本原理 2.1 插入Page 插入Page主要由读Page操作触发,上层函数读取Page的主要流程如下: 上层模块统一调用ReadBuffer()接口,读取指定表的指定Page ReadBuffer()中,先从Buffer查找Page,如果找到Page,则直接返回Buffer中的Page 如果未从Buffer找到Page,现在Buffer中预留存储Page的空间 之后,调用Smgr模块,从磁盘读取指定Page,并将其放入预留的空间内,再返回Buffer中的Page 函数调用关系大概如下图所示: ReadBuffer(Relation, BlockNumber) # 1 ReadBufferExtended() ReadBuffer_common() BufferAlloc() # 2,3 查找Page,如找到,返回;如未找到,预留存储Page的空间 smgrread(.., BlockNumber) # 4 如果未找到Page,从磁盘将Page读入预留的Page空间内,返回 2.2 查找Page 上层函数读Page时,首先进入查找Page阶段。 1.1 buffer池模型 模型: 1.2 buffer池接口 缓冲区接口: /* bufmgr.h */ /* 1 buffer池管理 */ void InitBufferPool(void); void InitBufferPoolAccess(void); void InitBufferPoolBackend(void); /* 2 buffer管理 */ /* 2.1 读buffer */ Buffer ReadBuffer(Relation reln, BlockNumber blockNum); Buffer ReadBufferExtended(Relation reln, ForkNumber forkNum, BlockNumber blockNum, ReadBufferMode mode, BufferAccessStrategy strategy); Buffer ReadBufferWithoutRelcache(RelFileNode rnode, ForkNumber forkNum, BlockNumber blockNum,ReadBufferMode mode, BufferAccessStrategy strategy); /* 2.2 清理buffer */ oid ReleaseBuffer(Buffer buffer); void UnlockReleaseBuffer(Buffer buffer); void MarkBufferDirty(Buffer buffer); void IncrBufferRefCount(Buffer buffer); Buffer ReleaseAndReadBuffer(Buffer buffer, Relation relation, BlockNumber blockNum); /* 2.3 刷盘buffer */ void CheckPointBuffers(int flags); void FlushOneBuffer(Buffer buffer); void FlushRelationBuffers(Relation rel); void FlushDatabaseBuffers(Oid dbid); void BufmgrCommit(void); bool BgBufferSync(void); /* 2.4 删除关联buffer */ void DropRelFileNodeBuffers(RelFileNodeBackend rnode, ForkNumber forkNum, BlockNumber firstDelBlock); void DropRelFileNodesAllBuffers(RelFileNodeBackend *rnodes, int nnodes); void DropDatabaseBuffers(Oid dbid); /* 2.5 锁buffer */ void UnlockBuffers(void); void LockBuffer(Buffer buffer, int mode); bool ConditionalLockBuffer(Buffer buffer); void LockBufferForCleanup(Buffer buffer); bool ConditionalLockBufferForCleanup(Buffer buffer); bool HoldingBufferPinThatDelaysRecovery(void); 1.3 buffer调用 插入数据 ...

January 21, 2025 · 3 min · 547 words · Me

PostgreSQL 3-7 Xlog

5 wal落盘 WalWriterMain for :: XLogBackgroundFlush LWLockAcquire(LW_EXCLUSIVE) XLogWrite LWLockRelease() CommitTransaction RecordTransactionCommit BufmgrCommit XactLogCommitRecord # xlog记录commit TransactionTreeSetCommitTsData XLogFlush(XactLastRecEnd) for :: WaitXLogInsertionsToFinish XLogWrite write ResourceOwnerRelease AtEOXact_Buffers AtEOXact_RelationCache ... 1 概述 1.1 背景 https://zhmin.github.io/posts/postgresql-wal-format/ wal机制解决的问题:持久性、原子性、性能、数据同步、恢复 持久性需求: 执行一次事务/事务块,例如INSERT语句,除了要写数据文件,还需要更新系统表、索引文件、fsm文件、vm文件、事务提交文件等 为避免故障导致数据丢失,产生的数据写入磁盘才算事务成功 应用 数据库 磁盘 | 1. INSERT (5) --> | | 2. 读取数据页 <-- | | 3. 将数据写入数据页 | 4. 存储数据页 --> | | 5. 读取索引文件 --> | | 6... | 7. 读取fsm文件 --> | | 8... | 9. 读取vm文件 --> | | 10... | 11. INSERT 成功 <-- | 原子性需求: ...

January 21, 2025 · 11 min · 2275 words · Me

PostgreSQL 3-2 Transation (v2)

概述 假如,2个会话,并发写入数据: BEGIN; INSERT INTO t1 VALUES(1, 'data1'); COMMIT; BEGIN; INSERT INTO t1 VALUES(1, 'data2'); COMMIT; 加锁逻辑: exec_simple_query pg_parse_query pg_analyze_and_rewrite pg_plan_queries PortalStart PortalRun PortalRunMulti ProcessQuery ExecutorStart standard_ExecutorStart InitPlan ExecInitNode ExecInitModifyTable ExecInitResultRelation ExecGetRangeTableRelation table_open(relation, 'RowExclusiveLock') # 8级锁,行级排他锁 relation_open LockRelationOid LockAcquire ExecutorRun standard_ExecutorRun ExecutePlan ... ExecInsert table_tuple_insert heapam_tuple_insert heap_insert RelationGetBufferForTuple LockBuffer LWLockAcquire(buffer, 'LW_EXCLUSIVE') # 读写锁 RelationPutHeapTuple log_heap_insert 锁 HWLock 8级 relation, index, buffer 死锁检测 LWLocks buffer SpinLock PresicateLock HWLock 锁 取值 操作 AccessShareLock 1 SELECT RowShareLock 2 SELECT FOR UPDATE/SHARE RowExclusiveLock 3 INSERT/UPDATE/DELETE ShareUpdateExclusiveLock 4 VACUUM/ANALYZE/CRATE INDEX CONCURRENTLY ShareLock 5 CREATE INDEX ShareRowExclusiveLock 6 row share ExclusiveLock 7 block ROW SHARE/SELECT FOR UPDATE AccessExclusiveLock 8 ALTER TABLE/DROP TABLE/VACUUM FULL LockRelationOid LockAcquireExtended hash_search 所有进程共享 ...

July 3, 2026 · 1 min · 166 words · Me

PostgreSQL 3-2 Transation

1 事务的概念 1.1 事务的基本特性 事务的4个特性: 原子性 Automicity 一致性 Consistency 隔离性 Isolation 持久性 Durability 1.2 事务的隔离性 不同隔离级别下,可能出现的异常如下:(T 表示可能) 异常 读未提交 读已提交 可重复读 可串行化 解释 脏读 T . . . 读其他事务未提交数据 不可重复读 T T . . 2次读取数据不一致 幻读 T T T . 2次读取数据条数不一致 postgresql默认的隔离级别是读已提交 -- 查看默认事务隔离级别 show default_transaction_isolation; default_transaction_isolation ------------------------------- read committed -- 查询当前的事务快照 pg 10 SELECT txid_current_snapshot(); txid_current_snapshot ----------------------- 1081:1081: -- 查询所有正在活跃的进程 SELECT * FROM pg_stat_activity; -- 1.3 事务的实现 为实现事务,有3种并发控制机制: ...

January 21, 2025 · 10 min · 1944 words · Me

PostgreSQL 3-5 FSM

1 概述 1.1 思想 fsm结构 用数组表示完全二叉树 | 8 | | 4 | 8 | | 4 | 2 | 1 | 8 | 1.2 FSM涉及 物理页号 +-------+-------+-------+-------+-------+ FSMFile | 8k | 8k | 8k | 8k | ... | +-------+-------+-------+-------+-------+ FSMBlock | 0 | 1 | 2 | 3 | ... | +-------+-------+-------+-------+-------+ 逻辑页号 第2层:1个FSMBlock 第1层:4069个FSMBlock 第0层:4069 * 4069个FSMBlock,每个FSM可映射4069个DataBlock FSMAddress: [层次,页号],页号即FSMBlock号 每个表最大容量是32T,只需3层也页面即可管理整个表的空闲空间 # [FSMAddress.level, FSMAddress.logpageno, FSMBlockNum(1个page即1个FSMBlock)] [2,0,0] | +---------------------------+---------------------------+ | | | [1,0,1] [1,1,1*4070+1] ... [1,4069,4069*4070+1] 本层共4069个FSMBlock | +----------+------------+ | | | [2,0,1+1] [2,1,1+2] ... [2,2,1+4069] ...... 本层共4069*4069个FSMBlock | (映射DataBlock 0 ~ 4065) 整体架构 3层FSMBlock [8] [8 0] [8 0 0 0] [...] [8 0 0 0 0 ...] [8] [8 0] [8 0 0 0] [...] [8 0 0 0 0 ...] [8] [8 0] [8 0 0 0] [...] [8 0 0 0 0 ...] 测试 ...

January 21, 2025 · 7 min · 1396 words · Me

PostgreSQL 3-10 VM

1 概述 1.1 vm接口调用点 heap_insert() buffer = RelationGetBufferForTuple() RelationPutHeapTuple(buffer, heaptup) if PageIsAllVisible(buffer): PageClearAllVisible(buffer) visibilitymap_clear(ItemPointerGetBlockNumber(heaptup), vmbuffer) MarkBufferDirty(buffer) #define PageIsAllVisible(page) ((page)->pd_flags & PD_ALL_VISIBLE) 1.2 vm模块对外接口 visibilitymap_clear(Relation rel, BlockNumber heapBlk, Buffer buf, uint8 flags) int mapByte = HEAPBLK_TO_MAPBYTE(heapBlk) LockBuffer(buf, BUFFER_LOCK_EXCLUSIVE) char *map = PageGetContents(buf) if map[mapByte] 1.3 vmmap结构与定位 /* 8192 - 24 = 8168 */ #define MAPSIZE (BLCKSZ - MAXALIGN(SizeOfPageHeaderData)) /* 1字节 2bit表示1页 : | 00 00 00 00 | */ #define BITS_PER_BYTE 8 #define BITS_PER_HEAPBLOCK 2 /* 1字节 可映射4个heapblock */ #define HEAPBLOCKS_PER_BYTE (BITS_PER_BYTE / BITS_PER_HEAPBLOCK) /* 计算heapblock对应的vmblock, vmbyte, vmoffset的位置 */ /* 1vmblock可映射4 * 8168个datablock */ #define HEAPBLOCKS_PER_PAGE (MAPSIZE * HEAPBLOCKS_PER_BYTE) #define HEAPBLK_TO_MAPBLOCK(x) ((x) / HEAPBLOCKS_PER_PAGE) #define HEAPBLK_TO_MAPBYTE(x) (((x) % HEAPBLOCKS_PER_PAGE) / HEAPBLOCKS_PER_BYTE) #define HEAPBLK_TO_OFFSET(x) (((x) % HEAPBLOCKS_PER_BYTE) * BITS_PER_HEAPBLOCK) #define VISIBILITYMAP_ALL_VISIBLE 0x01 #define VISIBILITYMAP_ALL_FROZEN 0x02 #define VISIBILITYMAP_VALID_BITS 0x03 1.4 vmblock读写 vm_extend(Relation rel, BlockNumber vm_nblocks) PGAlignedBlock pg PageInit((Page) pg.data, BLCKSZ, 0) LockRelationForExtension(rel, ExclusiveLock) BlockNumber vm_nblocks_now = smgrnblocks(rel->rd_smgr, VISIBILITYMAP_FORKNUM) while vm_nblocks_now < vm_nblocks: PageSetChecksumInplace((Page) pg.data, vm_nblocks_now) smgrextend(rel->rd_smgr, VISIBILITYMAP_FORKNUM, vm_nblocks_now, pg.data, false) vm_nblocks_now++ /* 让其他postgres进程smgrclose(rel) */ CacheInvalidateSmgr(rel->rd_smgr->smgr_rnode) { SharedInvalidationMessage msg msg.sm.id = SHAREDINVALSMGR_ID; msg.sm.backend_hi = rnode.backend >> 16; msg.sm.backend_lo = rnode.backend & 0xffff; msg.sm.rnode = rnode.node; VALGRIND_MAKE_MEM_DEFINED(&msg, sizeof(msg)) SendSharedInvalidMessages(&msg, 1) } rel->rd_smgr->smgr_vm_nblocks = vm_nblocks_now UnlockRelationForExtension(rel, ExclusiveLock) Buffer vm_readbuf(Relation rel, BlockNumber blkno, bool extend) RelationOpenSmgr(rel) if blkno >= rel->rd_smgr->smgr_vm_nblocks: if extend: vm_extend(rel, blkno + 1) buf = ReadBufferExtended(rel , VISIBILITYMAP_FORKNUM, blkno,RBM_ZERO_ON_ERROR) if PageIsNew(buf): PageInit(buf, BLKSZ, 0) return buf 1.5 vm模块对外接口 void visibilitymap_set(Relation rel, BlockNumber heapBlk, Buffer heapBuf, XLogRecPtr recptr, Buffer vmBuf, TransactionId cutoff_xid, uint8 flags) uint8 visibilitymap_get_status(Relation rel, BlockNumber heapBlk, Buffer *vmbuf) bool visibilitymap_clear(Relation rel, BlockNumber heapBlk, Buffer vmbuf, uint8 flags) /* 每个heapblock有3种状态 */ uint8 visibilitymap_get_status(Relation rel, BlockNumber heapBlk, Buffer *buf) BlockNumber mapBlock = HEAPBLK_TO_MAPBLOCK(heapBlk) uint32 mapByte = HEAPBLK_TO_MAPBYTE(heapBlk) uint8 mapOffset = HEAPBLK_TO_OFFSET(heapBlk) /* 比较传入的buf和计算的mapblock是否为同一个 */ if BufferIsValid(*buf): if BufferGetBlockNumber(*buf) != mapBlock: *buf = InvalidBuffer if !BufferIsValid(*buf): *buf = vm_readbuf(rel, mapBlock, false) map = PageGetContents(BufferGetPage(*buf)) /* 每个heapbuf有01 10 11共3种状态 */ unit8 result = ((map[mapByte] >> mapOffset) & VISIBILITYMAP_VALID_BITS) return uint8 void visibilitymap_set(Relation rel, BlockNumber heapBlk, Buffer heapBuf, XLogRecPtr recptr, Buffer vmBuf, TransactionId cutoff_xid, uint8 flags) ...

January 21, 2025 · 2 min · 330 words · Me

PostgreSQL 3-1 Heapam

1 heamp的简介 前置知识:Relation的基本概念,Tuple的基本概念 1.1 heapam接口的功能 heapam是存储模块对上层提供的数据读写接口,这些接口通常是操作1个relation中的1个或多个tuple。 本文介绍以下4个最常用的接口,即如何向1个relation中写入、删除、更新、读取1个Tuple,接口定义如下: /* 向relatiion中写入1条tuple */ heap_insert(Relation relation, HeapTuple tup, ...) /* 从relation中删除1条tuple */ heap_delete(Relation relation, ItemPointer tid, ...) /* 在relation中,将旧tuple标记删除,并写入1个新tuple */ heap_update(Relation relation, ItemPointer otid, HeapTuple newtup, ...) /* 从relation中读取数据,每次读取1条数据 */ HeapScanDesc heap_beginscan(Relation relation, ...) HeapTuple heap_getnext(HeapScanDesc scan) 1.2 调用heapam接口 /* * 此处,以4个常见的SQL为例,介绍数据库如何调用heapam接口 * 1. INSERT INTO t1 VALUES (1, 'data1'); * 2. SELECT * FROM t1; * 3. DELETE FROM t1 WHERE c1 = 1; * 4. UPDATE t1 SET c2 = 'data2' WHERE c1 = 1; */ exec_simple_query("SQL语句") /* 该函数是所有SQL语句的统一处理入口 */ PortalRun() PortalRunMulti() ProcessQuery() ExecutorRun() standard_ExecutorRun() ExecutePlan() ExecProcNode() /* 从这里开始分叉,不同类型的SQL语句,由不同函数处理 */ ExecModifyTable() ExecInsert() /* 处理:INSERT INTO t1 VALUES (1, 'data1') */ heap_insert(relation, tuple, ..) 1.3 heapam接口的使用场景 实际应用场景中,数据库的存储模块需满足以下关键需求: ...

January 21, 2025 · 3 min · 449 words · Me

PostgreSQL 3-8 Index

1 背景 1.1 数据检索 1 索引概述 1.1 概述 1.2 索引函数 1.2.1 调用索引函数 1.2.2 索引函数 2 索引实现 2.1 创建索引 ambuild 一、build整体流程 二、build插入tuple 2.2 索引写入数据 aminsert 3 operate 1 背景 1.1 数据检索 通常,应用会使用数据库存储大量数据,比如,单个表存储上百万、甚至上百亿条数据。此时,对表中数据进行某些查询时,将变得非常困难,比如等值查询、范围查询、模糊查询等。 1 索引概述 1.1 概述 索引方式:postgresql共有5种索引 唯一索引:不能出现重复的值 主键索引:主键自动创建唯一索引 多属性索引:多个列(最多32列) 部分索引:WHERE过滤索引 -- 示例 CREATE INDEX .. WHERE (c1 > 10); 表达式索引:和部分索引没太大区别,都是额外调用1个函数 -- 示例 CREATE INDEX .. WHERE (c1 = 10); 索引分类:postgresql有多种索引,本文见介绍以下4种 btree hash gist gin 索引系统表 pg_am:存储所有索引的方法 SELECT * FROM pg_am; amname | amhandler | amtype -------+-------------+-------- btree | bthandler | i hash | hashhandler | i gist | gisthandler | i gin | ginhandler | i spgist | spghandler | i brin | brinhandler | i pg_index:存储所有索引,索引列的下标 ...

January 21, 2025 · 7 min · 1391 words · Me

PostgreSQL 3-9 Vacuum

1 概述 select reltuples,relhasoids,oid from pg_class where relname = 't1'; -- reltuples = 0 select pg_relation_filepath('t1'); -- oid = filepath = relfilenode vacuum full t1; select reltuples,relhasoids,oid from pg_class where relname = 't1'; -- oid 不变, reltuples = 0 -- oid != filepath = relfilenode != analyze t1; -- reltuples = 2,即有效tuple 启动autovacuum进程 standard_ProcessUtility(Node *parsetree) case T_VacuumStmt: ExecVacuum(VacuumStmt *vacstmt) vacuum(RangeVar *relation, Oid relid, VacuumParams *params) ServerLoop(void) StartAutoVacLauncher(void) AutoVacLauncherMain() /* if fail: 发送信号重启aotuvacuum进程 */ SendPostmasterSignal(PMSIGNAL_START_AUTOVAC_WORKER) sigusr1_handler(SIGNAL_ARGS) if CheckPostmasterSignal(PMSIGNAL_START_AUTOVAC_WORKER): StartAutovacuumWorker(void) Backend *bn = malloc(); bn->pid = StartAutoVacWorker(void) AutoVacWorkerMain(int argc, char *argv[])() BaseInit(); InitPostgres(NULL, dbid, NULL, InvalidOid, dbname) do_autovacuum() autovacuum_do_vac_analyze(autovac_table *tab) vacuum(RangeVar *relation, Oid relid, VacuumParams *params) proc_exit(0) dlist_push_head(&BackendList, &bn->elem) 查找所有待vacuum的表 do_autovacuum(void) /* table-by-table */ StartTransactionCommand() /* 查询pg_database系统表 */ tuple = SearchSysCache1(DATABASEOID) /* 打开pg_class系统表 */ classRel = heap_open(RelationRelationId) pg_class_desc = CreateTupleDescCopy(RelationGetDescr(classRel)) /* 为toast表创建hash表 */ table_toast_map = hash_create() /* 遍历pg_class系统表,记录所有需要vacuum和analyze的普通表和系统表 */ relScan = heap_beginscan_catalog(classRel) while tuple = heap_getnext(relScan) != NULL: Form_pg_class classForm = (Form_pg_class) GETSTRUCT(tuple) relid = HeapTupleGetOid(tuple) relopts = extract_autovac_opts(tuple, pg_class_desc) tabentry = get_pgstat_tabentry_relid(relid, classForm->relisshared) /* 检查表是否需要被vacuum和analyze */ relation_needs_vacanalyze(relid, relopts, classForm, tabentry, &dovacuum, &doanalyze) if dovacuum || doanalyze: /* 普通表加入list中 */ table_oids = lappend_oid(table_oids, relid) /* toast表加入hash表中 */ if OidIsValid(classForm->reltoastrelid): hentry = hash_search(table_toast_map, &classForm->reltoastrelid, HASH_ENTER) hentry->ar_relid = relid eap_endscan(relScan) /* 第二次遍历pg_class系统表 */ relScan = heap_beginscan_catalog(classRel) while tuple = heap_getnext(relScan) != NULL: Form_pg_class classForm = (Form_pg_class) GETSTRUCT(tuple) relid = HeapTupleGetOid(tuple) relopts = extract_autovac_opts(tuple, pg_class_desc) ... bstrategy = GetAccessStrategy(BAS_VACUUM) /* 开始遍历所有list中的表 */ foreach(cell, table_oids) Oid relid = lfirst_oid(cell) /* 再次检查是否仍需vacuum */ tab = table_recheck_autovac(relid, table_toast_map, pg_class_desc) if tab == NULL: continue MyWorkerInfo->wi_tableoid = relid autovac_balance_cost() AutoVacuumUpdateDelay() autovacuum_do_vac_analyze(tab, bstrategy) vac_update_datfrozenxid() CommitTransactionCommand() 对单个表进行vacuum ...

January 21, 2025 · 4 min · 657 words · Me

PostgreSQL 3-13 逻辑复制测试记录

1 本地自测 1.1 快速tpcc vkill rm -rf $GAUSSHOME/data* # cp -r ~/w100_t200/data $GAUSSHOME/data # cp -r ~/w100_t200/data1 $GAUSSHOME/ cp -r ~/w5/data $GAUSSHOME/data cp -r ~/w5/data1 $GAUSSHOME/ # cp -r ~/w100/data $GAUSSHOME/data # cp -r ~/w100/data1 $GAUSSHOME/ # vb_guc set -D $GAUSSHOME/data -c "enable_ddl_logical_decode=on" # vb_guc set -D $GAUSSHOME/data1 -c "enable_ddl_logical_decode=on" # vb_guc set -D $GAUSSHOME/data1 -c "max_replication_slots=50" # vb_guc set -D $GAUSSHOME/data1 -c "max_background_workers=50" # vb_guc set -D $GAUSSHOME/data1 -c "max_logical_replication_workers=50" vb_guc set -D $GAUSSHOME/data1 -c "max_connections=4000" vb_guc set -D $GAUSSHOME/data1 -c "enable_audit=on" vb_guc set -D $GAUSSHOME/data1 -c "audit_buffer_size=262144" vb_guc set -D $GAUSSHOME/data1 -c "audit_directory_size=1073741824" vb_ctl start -D $GAUSSHOME/data vb_ctl start -D $GAUSSHOME/data1 # 五、发布端(发布) vsql -d pubdb -p 64000 -c "CREATE PUBLICATION pub1 FOR all tables WITH (ddl='all');" vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" # # 六、订阅端(订阅) gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription vsql -d pubdb -p 65000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=220);" # # vsql -d pubdb -p 65000 -r # \q # cd ~/tpcc # source envbm # ./toolbm run vpub # 换窗口查询 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vsql -d pubdb -p 64000 -c "SELECT * FROM vb_parallel_apply_stat()" vsql -d pubdb -p 64000 -c "SELECT * FROM vb_parallel_apply_stat()" # 场景一:验证内存 vsql -d pubdb -p 64000 -c "SELECT * FROM gs_thread_memory_context ORDER BY totalsize DESC limit 30;" > start_pub.mem vsql -d pubdb -p 64000 -c "SELECT * FROM gs_shared_memory_detail ORDER BY totalsize DESC limit 30;;" >> start_pub.mem vsql -d pubdb -p 65000 -c "SELECT * FROM gs_thread_memory_context ORDER BY totalsize DESC limit 30;;" > start_sub.mem vsql -d pubdb -p 65000 -c "SELECT * FROM gs_shared_memory_detail ORDER BY totalsize DESC limit 30;;" >> start_sub.mem # 场景八:验证数据一致性 vsql -d pubdb -p 64000 -c "DROP TABLE IF EXISTS tcnt" vsql -d pubdb -p 64000 -c "CREATE TABLE tcnt (c1 TEXT,c2 INT);" vsql -d pubdb -p 65000 -c "DROP TABLE IF EXISTS tcnt" vsql -d pubdb -p 65000 -c "CREATE TABLE tcnt (c1 TEXT,c2 INT);" for t in bmsql_config bmsql_customer bmsql_district bmsql_history bmsql_item bmsql_new_order bmsql_oorder bmsql_order_line bmsql_stock bmsql_warehouse; do echo "$t" vsql -d pubdb -p 64000 -c "INSERT INTO tcnt VALUES('$t', (SELECT count(1) FROM $t));" done; for t in bmsql_config bmsql_customer bmsql_district bmsql_history bmsql_item bmsql_new_order bmsql_oorder bmsql_order_line bmsql_stock bmsql_warehouse; do echo "$t" vsql -d pubdb -p 65000 -c "INSERT INTO tcnt VALUES('$t', (SELECT count(1) FROM $t));" done; vsql -d pubdb -p 64000 -c "SELECT * FROM tcnt" vsql -d pubdb -p 65000 -c "SELECT * FROM tcnt;" # 场景一:重启订阅端(kill) for i in $(seq 0 10); do echo "$i------------------------------------------------------------------------------------" ps ux | grep 'vastbase.*-D.*/data1' | grep -v grep | awk '{print $2}' | xargs kill -9 vb_ctl start -D $GAUSSHOME/data1 sleep 7 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" >> a.out done; # 场景二:重启订阅端(stop) for i in $(seq 0 10); do echo "$i------------------------------------------------------------------------------------" echo "$i------------------------------------------------------------------------------------" >> a.out vb_ctl stop -D $GAUSSHOME/data1 vb_ctl start -D $GAUSSHOME/data1 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" >> a.out sleep 7 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" >> a.out done; # 场景三:重启发布端 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vb_ctl stop -D $GAUSSHOME/data sleep 5 vb_ctl start -D $GAUSSHOME/data vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vb_ctl restart -D $GAUSSHOME/data1 # 场景三:重建 subscription for i in $(seq 0 10); do echo "$i------------------------------------------------------------------------------------" echo "$i------------------------------------------------------------------------------------" >> a.out vsql -d pubdb -p 65000 -c "DROP SUBSCRIPTION sub1;" vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" vsql -d pubdb -p 65000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=20);" sleep 6 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" >> a.out done; vsql -d pubdb -p 65000 -c "DROP SUBSCRIPTION sub1;" vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" vsql -d pubdb -p 65000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=20);" vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vsql -d pubdb -p 65000 -c "SELECT * FROM pg_subscription_rel" vsql -d pubdb -p 64000 -c "SELECT * FROM vb_parallel_apply_stat()" SELECT pg_wal_lsn_diff(pg_current_wal_lsn(), '0/0') / (1024 * 1024); # 场景五:ddl vsql -d pubdb -p 64000 -c "\d" vsql -d pubdb -p 65000 -c "\d" vsql -d pubdb -p 64000 -c "DROP TABLE bmsql_warehouse" vsql -d pubdb -p 64000 -c "CREATE TABLE t3(c1 INT)" vsql -d pubdb -p 64000 -c "\d" vsql -d pubdb -p 65000 -c "\d" # 场景六:多个发布订阅 vsql -d postgres -p 64000 -c "CREATE DATABASE db3;" vsql -d db3 -p 64000 -c "CREATE TABLE t3(c1 INT PRIMARY KEY, c2 TEXT);" vsql -d db3 -p 64000 -c "CREATE PUBLICATION pub3 FOR all tables WITH (ddl='all');" vsql -d db3 -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot3', 'pgoutput');" vsql -d postgres -p 65000 -c "CREATE DATABASE db3;" vsql -d db3 -p 65000 -c "CREATE TABLE t3(c1 INT PRIMARY KEY, c2 TEXT);" vsql -d db3 -p 65000 -c "CREATE SUBSCRIPTION sub3 CONNECTION 'host=172.16.103.90 port=64001 dbname=db3 user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub3 WITH (create_slot=false, slot_name = 'pslot3', copy_data = false, worker_number=20);" vsql -d pubdb -p 65000 -c "DROP SUBSCRIPTION sub1;" for i in $(seq 0 10); do vsql -d db3 -p 64000 -c "BEGIN; INSERT INTO t3 VALUES($i,'a'); UPDATE t3 SET c2='b' WHERE c1=$i-1; COMMIT;" & done; wait # 跑10分钟 ./toolbm run vpub vsql -d pubdb -p 65000 -c "DROP SUBSCRIPTION sub1;" # for t in bmsql_config bmsql_customer bmsql_district bmsql_history bmsql_item bmsql_new_order bmsql_oorder bmsql_order_line bmsql_stock bmsql_warehouse; do echo "$t" vsql -d pubdb -p 65000 -c "TRUNCATE $t" done; 续点重传 vkill rm -rf $GAUSSHOME/data* # cp -r ~/w5/data $GAUSSHOME/data # cp -r ~/w5/data1 $GAUSSHOME/ cp -r ~/w100/data $GAUSSHOME/data cp -r ~/w100/data1 $GAUSSHOME/ vb_guc set -D $GAUSSHOME/data1 -c "max_connections=4000" vb_ctl start -D $GAUSSHOME/data vb_ctl start -D $GAUSSHOME/data1 # 五、发布端(发布) vsql -d pubdb -p 64000 -c "CREATE PUBLICATION pub1 FOR all tables WITH (ddl='all');" vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" # # 六、订阅端(订阅) gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription vsql -d pubdb -p 65000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=220);" # # vsql -d pubdb -p 65000 -r # \q vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vsql -d pubdb -p 64000 -c "SELECT * FROM vb_parallel_apply_stat()" # ./toolbm run vpub # 重启订阅端(kill) for i in $(seq 0 10); do echo "$i------------------------------------------------------------------------------------" ps ux | grep 'vastbase.*-D.*/data1' | grep -v grep | awk '{print $2}' | xargs kill -9 vb_ctl start -D $GAUSSHOME/data1 sleep 7 vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" >> a.out done; vsql -d pubdb -p 64000 -c " SELECT * FROM pg_get_replication_slots(); " vsql -d pubdb -p 64000 -c " SELECT pg_current_wal_lsn();" # 验证数据一致性 vsql -d pubdb -p 64000 -c "DROP TABLE IF EXISTS tcnt" vsql -d pubdb -p 64000 -c "CREATE TABLE tcnt (c1 TEXT,c2 INT);" vsql -d pubdb -p 65000 -c "DROP TABLE IF EXISTS tcnt" vsql -d pubdb -p 65000 -c "CREATE TABLE tcnt (c1 TEXT,c2 INT);" for t in bmsql_config bmsql_customer bmsql_district bmsql_history bmsql_item bmsql_new_order bmsql_oorder bmsql_order_line bmsql_stock bmsql_warehouse; do echo "$t" vsql -d pubdb -p 64000 -c "INSERT INTO tcnt VALUES('$t', (SELECT count(1) FROM $t));" done; for t in bmsql_config bmsql_customer bmsql_district bmsql_history bmsql_item bmsql_new_order bmsql_oorder bmsql_order_line bmsql_stock bmsql_warehouse; do echo "$t" vsql -d pubdb -p 65000 -c "INSERT INTO tcnt VALUES('$t', (SELECT count(1) FROM $t));" done; vsql -d pubdb -p 64000 -c "SELECT * FROM tcnt" vsql -d pubdb -p 65000 -c "SELECT * FROM tcnt;" -- ===================== worker_number ======================= CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false); DROP SUBSCRIPTION sub1; \! vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" CREATE SUBSCRIPTION sub2 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=1); DROP SUBSCRIPTION sub2; \! vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" CREATE SUBSCRIPTION sub3 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=2); DROP SUBSCRIPTION sub3; \! vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" CREATE SUBSCRIPTION sub4 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=200); DROP SUBSCRIPTION sub4; \! vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" CREATE SUBSCRIPTION sub5 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=400); DROP SUBSCRIPTION sub5; CREATE SUBSCRIPTION sub6 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=1000); DROP SUBSCRIPTION sub6; -- ===================== copy_data ======================= -- do truncate fist CREATE SUBSCRIPTION sub21 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = true, worker_number=50); SELECT * FROM pg_subscription_rel; DROP SUBSCRIPTION sub21; -- copy and stop -- ===================== cache file ======================= -- 1. drop subscription -- 2. kill progress -- ===================== 断点续传 ======================= vsql -d postgres -p 50000 -c "DROP TABLE IF EXISTS t1" vsql -d postgres -p 50000 -c "CREATE TABLE t1(c1 INT, c2 TEXT)" for i in $(seq 0 10); do vsql -d postgres -p 50000 -c "BEGIN; INSERT INTO t1 VALUES($i,'a'); UPDATE t1 SET c2='b' WHERE c1=$i-1; COMMIT;" & done vsql -d postgres -p 50000 -c "SELECT * FROM t1 ORDER BY c1" parallel_apply.source 1 自测 批量验证 初始化模板 export VDATA=~/g100/rep/data vb_guc set -D $VDATA -c "port=64000" vb_guc set -D $VDATA -c "max_connection=2000" vb_guc set -D $VDATA -c "wal_level=logical" vb_guc set -D $VDATA -c "max_replication_slots=100" vb_guc set -D $VDATA -c "max_background_workers=100" vb_guc set -D $VDATA -c "max_logical_replication_workers=100" sed -i '1i host all all 0.0.0.0/0 sha256\n' $VDATA/pg_hba.conf echo "host replication all 0.0.0.0/0 sha256" >> $VDATA/pg_hba.conf insert性能 export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vstart vsql -d postgres -p 64000 -r -- 3 create publicatioin \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables; CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput'); -- sub CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name='pslot1', copy_data=false, worker_number=80); SELECT count(1) FROM t1; SELECT * FROM vb_parallel_apply_stat(); \c pubdb truncate t1; \c subdb truncate t1; SELECT datname,state,length(query) FROM pg_stat_activity WHERE datname = 'pubdb'; 多个发布订阅 vb_guc set -D ~/g100/rep/data -c "max_replication_slots=256" vb_guc set -D ~/g100/rep/data -c "max_background_workers=256" vb_guc set -D ~/g100/rep/data -c "max_logical_replication_workers=256" export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vb_guc set -D ~/g100/rep/data -c "enable_ddl_logical_decode=on" vb_guc set -D ~/g100/rep/data -c "max_wal_senders=36" vb_guc set -D ~/g100/rep/data -c "max_replication_slots=256" vb_guc set -D ~/g100/rep/data -c "max_background_workers=256" vb_guc set -D ~/g100/rep/data -c "max_logical_replication_workers=256" vb_guc set -D ~/g100/rep/data -c "max_connections=1000" vstart rm -rf ~/g100/rep/data1 vb_guc set -D ~/g100/rep/data1 -c "port=65000" cp -r ~/g100/rep/temp ~/g100/rep/data1 vb_ctl start -D ~/g100/rep/data1 # 发布 vsql -d postgres -p 64000 -c "CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345';" for i in $(seq 0 11); do vsql -d postgres -p 64000 -c "CREATE DATABASE pubdb$i;" vsql -d pubdb$i -p 64000 -c "CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT);" vsql -d pubdb$i -p 64000 -c "CREATE PUBLICATION pub$i FOR all tables;" vsql -d pubdb$i -p 64000 -c "SELECT pg_create_logical_replication_slot('pslot$i', 'pgoutput');" done; # 5 load data for i in $(seq 0 11); do for j in $(seq 0 10); do vsql -d pubdb$i -p 64000 -c "BEGIN; INSERT INTO t1 VALUES($j,'a'); UPDATE t1 SET c2='b' WHERE c1=$j-1; COMMIT;" & done; sleep 2 done; wait # 订阅 gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription for i in $(seq 0 11); do vsql -d postgres -p 65000 -c "CREATE DATABASE subdb$i;" vsql -d subdb$i -p 65000 -c "CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT);" vsql -d subdb$i -p 65000 -c "CREATE SUBSCRIPTION sub$i CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb$i user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub$i WITH (create_slot=false, slot_name = 'pslot$i', copy_data = false, worker_number=30);" done; sleep 5 # 6 check for i in $(seq 0 11); do vsql -d pubdb$i -p 64000 -c "SELECT count(1) FROM t1" done; for i in $(seq 0 11); do vsql -d subdb$i -p 65000 -c "SELECT count(1) FROM t1" done; vsql -d postgres -p 64000 -c " SELECT * FROM pg_get_replication_slots();" vsql -d postgres -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" # 7 cleanup DROP SUBSCRIPTION sub1; \c postgres DROP DATABASE IF EXISTS pubdb; DROP DATABASE IF EXISTS subdb; DROP USER pu1; copydata export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vstart rm -rf ~/g100/rep/data1 cp -r ~/g100/rep/temp ~/g100/rep/data1 vb_guc set -D ~/g100/rep/data1 -c "port=65000" vb_ctl start -D ~/g100/rep/data1 vsql -d postgres -p 64000 -r -- 3 create publicatioin -- pub \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables; -- CREATE PUBLICATION pub1 FOR TABLE t1,t2; CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; -- INSERT INTO t1 VALUES(1,'a'),(2,'b'),(3,'c'); -- INSERT INTO t1 (c1,c2) SELECT i,'BASE_' || i || '_' || repeat('X', 13000) FROM generate_series(1, 10000) AS i; \q -- for i in $(seq 0 1000); do vsql -d pubdb -p 64000 -c "INSERT INTO t1 VALUES($i,'a'); " ; done vsql -d postgres -p 65000 -r -- sub CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (copy_data=false, worker_number=10); -- \! for i in $(seq 20000 20010); do vsql -d pubdb -p 64000 -c "INSERT INTO t1 VALUES($i,'a'); " ; done SELECT * FROM t1; SELECT * FROM vb_parallel_apply_stat(); vsql -d pubdb -p 64000 -c "SELECT count(1) FROM t1" vsql -d subdb -p 65000 -c "SELECT count(1) FROM t1" vsql -d subdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" SELECT * FROM pg_replication_slots; INSERT INTO t1 VALUES(4,'a'),(5,'b'),(6,'c'); ddl export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vb_guc set -D ~/g100/rep/data -c "enable_ddl_logical_decode=on" vstart rm -rf ~/g100/rep/data1 cp -r ~/g100/rep/temp ~/g100/rep/data1 vb_guc set -D ~/g100/rep/data1 -c "port=65000" vb_ctl start -D ~/g100/rep/data1 vsql -d postgres -p 64000 -r -- 3 create publicatioin -- pub \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables WITH (publish='insert,update,delete,truncate', ddl='all'); -- CREATE PUBLICATION pub1 FOR TABLE t1,t2; CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; -- INSERT INTO t1 VALUES(1,'a'),(2,'b'),(3,'c'); \q vsql -d postgres -p 65000 -r -- sub CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (copy_data=true, worker_number=4); SELECT * FROM t1; \! vsql -d pubdb -p 64000 -c "TRUNCATE t1" \! sleep 1 SELECT * FROM t1; \! vsql -d pubdb -p 64000 -c "DROP TABLE t1" \! sleep 1 SELECT * FROM t1; SELECT * FROM pg_replication_slots; INSERT INTO t1 VALUES(4,'a'),(5,'b'),(6,'c'); sub-txn + ddl export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vb_guc set -D ~/g100/rep/data -c "enable_ddl_logical_decode=on" vstart rm -rf ~/g100/rep/data1 cp -r ~/g100/rep/temp ~/g100/rep/data1 vb_guc set -D ~/g100/rep/data1 -c "port=65000" vb_ctl start -D ~/g100/rep/data1 vsql -d postgres -p 64000 -r -- 3 create publicatioin -- pub \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables WITH (publish='insert,update,delete,truncate', ddl='all'); -- CREATE PUBLICATION pub1 FOR TABLE t1,t2; CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; INSERT INTO t1 VALUES(1,'a'),(2,'b'),(3,'c'); \q vsql -d postgres -p 65000 -r -- sub CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (copy_data=true, worker_number=4); SELECT * FROM t1; \! vsql -d pubdb -p 64000 -c "TRUNCATE t1" \! sleep 1 SELECT * FROM t1; \! vsql -d pubdb -p 64000 -c "DROP TABLE t1" \! sleep 1 SELECT * FROM t1; \! vsql -d pubdb -p 64000 -c "BEGIN; CREATE TABLE t3(c1 int); SAVEPOINT sp1; INSERT INTO t3 VALUES(3); SAVEPOINT sp2; TRUNCATE t3; SAVEPOINT sp3; COMMIT;" SELECT * FROM t3; DROP TABLE t3; SELECT * FROM pg_replication_slots; INSERT INTO t1 VALUES(4,'a'),(5,'b'),(6,'c'); disable export VDATA=~/g100/rep/data vkill rm -rf ~/g100/rep/data cp -r ~/g100/rep/temp ~/g100/rep/data vstart rm -rf ~/g100/rep/data1 cp -r ~/g100/rep/temp ~/g100/rep/data1 vb_guc set -D ~/g100/rep/data1 -c "port=65000" vb_ctl start -D ~/g100/rep/data1 vsql -d postgres -p 64000 -r -- 3 create publicatioin -- pub \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables WITH (ddl='all'); -- CREATE PUBLICATION pub1 FOR TABLE t1,t2; CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; INSERT INTO t1 VALUES(1,'a'),(2,'b'),(3,'c'); \q vsql -d postgres -p 65000 -r -- sub CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (copy_data=true, worker_number=4); SELECT * FROM t1; -- SELECT * FROM pg_replication_slots; SELECT * FROM pg_stat_subscription; ALTER SUBSCRIPTION sub1 DISABLE; \! vsql -d pubdb -p 64000 -c "INSERT INTO t1 VALUES(4,'a'),(5,'b'),(6,'c');" SELECT * FROM t1; ALTER SUBSCRIPTION sub1 ENABLE; SELECT * FROM pg_replication_slots; 基础场景 vkill vinit -- 1 config vb_guc set -D $GAUSSHOME/data -c "wal_level=logical" vb_guc set -D $GAUSSHOME/data -c "max_replication_slots=100" vb_guc set -D $GAUSSHOME/data -c "max_background_workers=100" vb_guc set -D $GAUSSHOME/data -c "max_logical_replication_workers=100" sed -i '1i host all all 0.0.0.0/0 sha256\n' $GAUSSHOME/data/pg_hba.conf echo "host replication all 0.0.0.0/0 sha256" >> $GAUSSHOME/data/pg_hba.conf -- 2 startup vstart vsql -d postgres -p 50000 -r -- 3 create publicatioin \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables; SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput'); CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; -- 4 create subscription CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription \! vsql -d subdb -p 50000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=50001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false);" -- , worker_number=30 -- 172.16.103.90 \c pubdb INSERT INTO t1 VALUES(1,'a'),(2,'b'),(3,'c'); UPDATE t1 SET c1 = 2 WHERE c2 = 'c'; INSERT INTO t1 VALUES(4,'d'); BEGIN; INSERT INTO t1 VALUES(5,'e'); UPDATE t1 SET c1 = 1 WHERE c2 = 'b'; INSERT INTO t1 VALUES(6,'f'); COMMIT; BEGIN; INSERT INTO t1 VALUES(7,'d'); UPDATE t1 SET c2 = 'a' WHERE c1 = 2; COMMIT; -- 5 load data \! for i in $(seq 0 10); do vsql -d pubdb -p 50000 -c "BEGIN; INSERT INTO t1 VALUES($i,'a'); UPDATE t1 SET c2='b' WHERE c1=$i-1; COMMIT;" > /dev/null 2>&1 & done; wait \! sleep 3 -- 5 check subscription \o pubdb.out \c pubdb SELECT * FROM t1 ORDER BY c1; \o \o subdb.out \c subdb SELECT * FROM t1 ORDER BY c1; \o \! grep -Irn "rows" pubdb.out \! grep -Irn "rows" subdb.out \! diff pubdb.out subdb.out \! cat subdb.out \! rm -f pubdb.out subdb.out -- 7 cleanup DROP SUBSCRIPTION sub1; \c postgres DROP DATABASE IF EXISTS pubdb; DROP DATABASE IF EXISTS subdb; DROP USER pu1; 空库feedback vkill rm -rf $GAUSSHOME/data* cp -r ~/w5/data $GAUSSHOME/data cp -r ~/w5/data1 $GAUSSHOME/ vb_guc set -D $GAUSSHOME/data1 -c "max_connections=4000" vb_ctl start -D $GAUSSHOME/data vb_ctl start -D $GAUSSHOME/data1 # 五、发布端(发布) vsql -d pubdb -p 64000 -c "CREATE PUBLICATION pub1 FOR all tables WITH (ddl='all');" vsql -d pubdb -p 64000 -c " SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput');" vsql -d pubdb -p 64000 -c "CREATE DATABASE db1"; vsql -d pubdb -p 64000 -c "CREATE DATABASE db2"; vsql -d db1 -p 64000 -c "CREATE PUBLICATION p1 FOR all tables WITH (ddl='all');" vsql -d db1 -p 64000 -c " SELECT pg_create_logical_replication_slot('ps1', 'pgoutput');" vsql -d db2 -p 64000 -c "CREATE PUBLICATION p2 FOR all tables WITH (ddl='all');" vsql -d db2 -p 64000 -c " SELECT pg_create_logical_replication_slot('ps2', 'pgoutput');" # 六、订阅端(订阅) gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription vsql -d pubdb -p 65000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=172.16.103.90 port=64001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=220);" vsql -d pubdb -p 65000 -c "CREATE DATABASE db1"; vsql -d pubdb -p 65000 -c "CREATE DATABASE db2"; vsql -d db1 -p 65000 -c "CREATE SUBSCRIPTION s1 CONNECTION 'host=172.16.103.90 port=64001 dbname=db1 user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION p1 WITH (create_slot=false, slot_name = 'ps1', copy_data = false, worker_number=220);" vsql -d db2 -p 65000 -c "CREATE SUBSCRIPTION s2 CONNECTION 'host=172.16.103.90 port=64001 dbname=db2 user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION p2 WITH (create_slot=false, slot_name = 'ps2', copy_data = false, worker_number=220);" vsql -d pubdb -p 64000 -c "CREATE VIEW pstat AS SELECT slot_name,active,confirmed_flush, (pg_wal_lsn_diff(confirmed_flush, '0/0') / (1024 * 1024))::int as confirmed_lsn FROM pg_get_replication_slots();" # 查看复制进度 vsql -d pubdb -p 64000 -c "SELECT * FROM pstat" # 查看速度 vsql -d pubdb -p 64000 -c "SELECT (pg_wal_lsn_diff(pg_current_wal_lsn(), '0/0') / (1024 * 1024))::int as mb" vsql -d pubdb -p 65000 -c "SELECT * FROM vb_parallel_apply_stat()" vsql -d pubdb -p 64000 -c "SELECT * FROM vb_parallel_apply_stat()" 子事务 1. 这是一个基于opengauss开发的数据 2. 我已经设置了一些alias命令。全量编译vcmake。增量编译:vrecmake。启动集群:vinit, vstart。强制停止集群:vkill。 3. core文件目录:/data/ 4. 并行逻辑复制关键代码:parallel_apply.cpp 5. $GAUSSHOME也就是install目录下,所有文件都是临时生成的,随便动 vkill vinit -- 1 config vb_guc set -D $GAUSSHOME/data -c "wal_level=logical" vb_guc set -D $GAUSSHOME/data -c "max_replication_slots=100" vb_guc set -D $GAUSSHOME/data -c "max_background_workers=100" vb_guc set -D $GAUSSHOME/data -c "max_logical_replication_workers=100" # vb_guc set -D $GAUSSHOME/data -c "enable_audit=on" # vb_guc set -D $GAUSSHOME/data -c "audit_buffer_size=262144" # vb_guc set -D $GAUSSHOME/data -c "audit_directory_size=1073741824" sed -i '1i host all all 0.0.0.0/0 sha256\n' $GAUSSHOME/data/pg_hba.conf echo "host replication all 0.0.0.0/0 sha256" >> $GAUSSHOME/data/pg_hba.conf -- 2 startup vstart vsql -d postgres -p 50000 -r -- 3 create publicatioin \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables; SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput'); CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; -- 4 create subscription CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription \! vsql -d subdb -p 50000 -c "CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=50001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=20);" -- , worker_number=30 -- 172.16.103.90 \c pubdb BEGIN; INSERT INTO t1 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t1 VALUES(2,'222'); -- 假如xid=2 [savepoint 2, xid:1-2] SAVEPOINT s2; INSERT INTO t1 VALUES(3,'333'); -- 假如xid=3 [savepoint 3, xid:1-2-3] SAVEPOINT s3; INSERT INTO t1 VALUES(4,'444'); -- 假如xid=4 [savepoint 4, xid:1-2-3-4] SAVEPOINT s4; INSERT INTO t1 VALUES(5,'555'); -- 假如xid=5 [savepoint 5, xid:1-2-3-4-5] ROLLBACK TO SAVEPOINT s2; INSERT INTO t1 VALUES(6,'666'); -- 假如xid=6 [rollback to 3: xid:1-2] SAVEPOINT s5; INSERT INTO t1 VALUES(7,'777'); -- 假如xid=7 RELEASE SAVEPOINT s2; -- 此时,xid=6和xid=7会归并到xid=2的事务中 INSERT INTO t1 VALUES(8,'888'); -- 假如xid=2(假设上面RELEASE到s1,则xid=1,和顶层事务一致) COMMIT; -- \! sleep 3 -- 5 check subscription \c pubdb SELECT xmin,* FROM t1 ORDER BY c1; \c subdb SELECT * FROM vb_parallel_apply_stat(); SELECT xmin,* FROM t1 ORDER BY c1; -- rollback to top \c pubdb BEGIN; INSERT INTO t2 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t2 VALUES(2,'222'); -- 假如xid=2 SAVEPOINT s2; INSERT INTO t2 VALUES(3,'333'); -- 假如xid=3 SAVEPOINT s3; INSERT INTO t2 VALUES(4,'444'); -- 假如xid=4 SAVEPOINT s4; INSERT INTO t2 VALUES(5,'555'); -- 假如xid=5 ROLLBACK TO SAVEPOINT s2; INSERT INTO t2 VALUES(6,'666'); -- 假如xid=6 SAVEPOINT s5; INSERT INTO t2 VALUES(7,'777'); -- 假如xid=7 RELEASE SAVEPOINT s1; -- 此时,xid=6和xid=7会归并到xid=2的事务中 INSERT INTO t2 VALUES(8,'888'); -- 假如xid=2(假设上面RELEASE到s1,则xid=1,和顶层事务一致) COMMIT; -- \! sleep 3 -- 5 check subscription \c pubdb SELECT xmin,* FROM t2 ORDER BY c1; \c subdb SELECT xmin,* FROM t2 ORDER BY c1; SELECT * FROM vb_parallel_apply_stat(); -- 延时检查 CREATE TABLE t3 (c1 INT, c2 TEXT ); ALTER TABLE t3 ADD CONSTRAINT t3_c1_unique UNIQUE (c1) DEFERRABLE; -- 事务1 -- 事务2 BEGIN; INSERT INTO t3 VALUES (1, 'a'); BEGIN; SET CONSTRAINTS ALL DEFERRED; INSERT INTO t3 VALUES (1, 'a'); COMMIT; COMMIT; -- 7 cleanup DROP SUBSCRIPTION sub1; \c postgres DROP DATABASE IF EXISTS pubdb; DROP DATABASE IF EXISTS subdb; DROP USER pu1; 子事务 + 主键冲突 vkill vinit -- 1 config vb_guc set -D $GAUSSHOME/data -c "wal_level=logical" vb_guc set -D $GAUSSHOME/data -c "max_replication_slots=100" vb_guc set -D $GAUSSHOME/data -c "max_background_workers=100" vb_guc set -D $GAUSSHOME/data -c "max_logical_replication_workers=100" sed -i '1i host all all 0.0.0.0/0 sha256\n' $GAUSSHOME/data/pg_hba.conf echo "host replication all 0.0.0.0/0 sha256" >> $GAUSSHOME/data/pg_hba.conf -- 2 startup vstart vsql -d postgres -p 50000 -r -- 3 create publicatioin \c postgres CREATE DATABASE pubdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR all tables; SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput'); CREATE USER pu1 REPLICATION SYSADMIN LOGIN ENCRYPTED PASSWORD 'pu1.12345'; -- 4 create subscription CREATE DATABASE subdb; \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); \! gs_guc generate -S 'pu1.12345' -D $GAUSSHOME/bin -o subscription CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=50001 dbname=pubdb user=pu1 password=pu1.12345 sslmode=disable' PUBLICATION pub1 WITH (create_slot=false, slot_name = 'pslot1', copy_data = false, worker_number=20); INSERT INTO t1 VALUES(3,'old'); \c pubdb BEGIN; INSERT INTO t1 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t1 VALUES(2,'222'); -- 假如xid=2 [savepoint 2, xid:1-2] SAVEPOINT s2; INSERT INTO t1 VALUES(3,'333'); -- 假如xid=3 [savepoint 3, xid:1-2-3] SAVEPOINT s3; INSERT INTO t1 VALUES(4,'444'); -- 假如xid=4 [savepoint 4, xid:1-2-3-4] SAVEPOINT s4; COMMIT; \c subdb SELECT * FROM t1; SELECT * FROM vb_parallel_apply_stat(); -- 7 cleanup DROP SUBSCRIPTION sub1; \c postgres DROP DATABASE IF EXISTS pubdb; DROP DATABASE IF EXISTS subdb; DROP USER pu1; 子事务开发 -- 子事务开发 CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); SELECT pg_current_wal_lsn(); -- 场景一:relase commit BEGIN; INSERT INTO t1 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t1 VALUES(2,'222'); -- 假如xid=2 SAVEPOINT s2; INSERT INTO t1 VALUES(3,'333'); -- 假如xid=3 SAVEPOINT s3; INSERT INTO t1 VALUES(4,'444'); -- 假如xid=4 SAVEPOINT s4; INSERT INTO t1 VALUES(5,'555'); -- 假如xid=5 ROLLBACK TO SAVEPOINT s2; INSERT INTO t1 VALUES(6,'666'); -- 假如xid=6 SAVEPOINT s5; INSERT INTO t1 VALUES(7,'777'); -- 假如xid=7 RELEASE SAVEPOINT s2; -- 此时,xid=6和xid=7会归并到xid=2的事务中 INSERT INTO t1 VALUES(8,'888'); -- 假如xid=2(假设上面RELEASE到s1,则xid=1,和顶层事务一致) COMMIT; SELECT xmin,* FROM t1 ORDER BY c1; \! pg_xlogdump -p $GAUSSHOME/data/pg_xlog -s xxx -x xid -- 场景二:release abort BEGIN; INSERT INTO t1 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t1 VALUES(2,'222'); -- 假如xid=2 SAVEPOINT s2; INSERT INTO t1 VALUES(3,'333'); -- 假如xid=3 SAVEPOINT s3; INSERT INTO t1 VALUES(4,'444'); -- 假如xid=4 SAVEPOINT s4; INSERT INTO t1 VALUES(5,'555'); -- 假如xid=5 ROLLBACK TO SAVEPOINT s2; INSERT INTO t1 VALUES(6,'666'); -- 假如xid=6 SAVEPOINT s5; INSERT INTO t1 VALUES(7,'777'); -- 假如xid=7 RELEASE SAVEPOINT s2; -- 此时,xid=6和xid=7会归并到xid=2的事务中 INSERT INTO t1 VALUES(8,'888'); -- 假如xid=2(假设上面RELEASE到s1,则xid=1,和顶层事务一致) ROLLBACK TO SAVEPOINT s1; COMMIT; SELECT xmin,* FROM t1 ORDER BY c1; -- pg 用例 echo " wal_level=logical max_worker_processes=20 max_logical_replication_workers=20 max_parallel_apply_workers_per_subscription=10 " >> $PG_HOME/data/postgresql.conf \c postgres CREATE DATABASE pubdb; CREATE DATABASE subdb; \c pubdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE PUBLICATION pub1 FOR ALL TABLES; CREATE USER pu1 REPLICATION LOGIN PASSWORD 'pu1.12345'; SELECT pg_create_logical_replication_slot('pslot1', 'pgoutput'); -- 2 create subscription \c subdb CREATE TABLE t1(c1 INT PRIMARY KEY, c2 TEXT); CREATE TABLE t2(c1 INT PRIMARY KEY, c2 TEXT); CREATE SUBSCRIPTION sub1 CONNECTION 'host=127.0.0.1 port=5432 dbname=pubdb user=pu1 password=pu1.12345 replication=database' PUBLICATION pub1 WITH ( create_slot=false, slot_name='pslot1', copy_data=false, streaming=parallel ); -- 3 test transaction -- \c pubdb -- BEGIN; -- INSERT INTO t1 VALUES(1,'111'); -- top xid -- SAVEPOINT s1; -- INSERT INTO t1 VALUES(2,'222'); -- subxid 2 -- SAVEPOINT s2; -- INSERT INTO t1 VALUES(3,'333'); -- subxid 3 -- SAVEPOINT s3; -- INSERT INTO t1 VALUES(4,'444'); -- subxid 4 -- SAVEPOINT s4; -- INSERT INTO t1 VALUES(5,'555'); -- subxid 5 -- ROLLBACK TO SAVEPOINT s2; -- INSERT INTO t1 VALUES(6,'666'); -- subxid 6 -- SAVEPOINT s5; -- INSERT INTO t1 VALUES(7,'777'); -- subxid 7 -- RELEASE SAVEPOINT s2; -- INSERT INTO t1 VALUES(8,'888'); -- 回到s1对应事务上下文 -- COMMIT; \c pubdb BEGIN; INSERT INTO t1 VALUES(1,'111'); -- 假如xid=1 SAVEPOINT s1; INSERT INTO t1 VALUES(2,'222'); -- 假如xid=2 SAVEPOINT s2; INSERT INTO t1 VALUES(3,'333'); -- 假如xid=3 SAVEPOINT s3; INSERT INTO t1 VALUES(4,'444'); -- 假如xid=4 SAVEPOINT s4; INSERT INTO t1 VALUES(5,'555'); -- 假如xid=5 ROLLBACK TO SAVEPOINT s2; INSERT INTO t1 VALUES(6,'666'); -- 假如xid=6 SAVEPOINT s5; INSERT INTO t1 VALUES(7,'777'); -- 假如xid=7 RELEASE SAVEPOINT s2; -- 此时,xid=6和xid=7会归并到xid=2的事务中 INSERT INTO t1 VALUES(8,'888'); -- 假如xid=2(假设上面RELEASE到s1,则xid=1,和顶层事务一致) ROLLBACK TO SAVEPOINT s1; COMMIT; -- 4 check \c pubdb SELECT xmin,* FROM t1 ORDER BY c1; \c subdb SELECT * FROM pg_stat_subscription; SELECT xmin,* FROM t1 ORDER BY c1; # 发布端:2阶段处理 ReleaseSavepoint for loop: cur = CurrentTransactionState cur->blockState = TBLOCK_SUBRELEASE cur = cur->parent CommitTransactionCommand if CurrentTransactionState is TBLOCK_SUBRELEASE for loop: CommitSubTransaction AtSubCommit_childXids CurrentTransactionState->parent->childXids[cnt++] = CurrentTransactionState->xid PopTransaction CurrentTransactionState = s->parent # 订阅端: pa_start_subtrans DefineSavepoint subxactlist = lappend_xid(subxactlist, current_xid) apply_handle_stream_commit pa_stream_abort for loop: curxid = subxactlist if curxid == subxid: RollbackToSavepoint 实现原理: ...

May 21, 2026 · 32 min · 6745 words · Me
心情不好的时候可以点一下 🐱
×
🤖 Doubao AI ×
Hi! 我是你的技术助手。关于代码、架构或 Bug,随时问我!🚀