首页 文章 精选 留言 我的

精选列表

搜索[文档解析],共10000篇文章
优秀的个人博客,低调大师

openGauss数据库源码解析系列文章——存储引擎源码解析(二)

上一篇我们讲述了“4.2 磁盘引擎”中“4.2.1 磁盘引擎整体框架及代码概览”与“4.2.2 行存储统一访存接口”。本篇我们将讲述“4.2.3 astore”。 4.2.3 astore astore整体框架 astore整体框架如图4-2所示。如上所述,作为行存储子格式之一,astore需要实现自己的堆表存取(访存)管理接口、堆表页面结构、堆表元组结构、元组多版本机制,以及空闲空间管理和回收机制。 图4-2 astore整体框架示意图 astore堆表页面元组结构 本节介绍astore堆表的页面和元组结构。 所谓堆表,是指元组无序存储,数据按照“先来后到”的方式存储在页面中的空闲位置。作为对比,在索引表中,元组根据索引键键值的排序,在页面内部有序存储,且各个页面之间在逻辑上也是有序存储的。堆表存储数据主体,索引表仅存储索引键键值以及对应的、完整元组的物理位置(即完整元组在堆表中的页面号和页内偏移)。 1) astore堆表元组结构 astore堆表元组结构的定义部分代码如下: typedef struct HeapTupleFields { ShortTransactionId t_xmin; /* 插入元组事务的事务号 */ ShortTransactionId t_xmax; /* 删除元组事务的事务号 */ union { CommandId t_cid; /* 插入或删除命令在事务中的命令号 */ ShortTransactionId t_xvac; } t_field3; } HeapTupleFields; typedef struct HeapTupleHeaderData { union { HeapTupleFields t_heap; DatumTupleFields t_datum; } t_choice; ItemPointerData t_ctid; /* 当前元组或更新后元组的行号 */ uint16 t_infomask2; /* 字段个数和标记位 */ uint16 t_infomask; /* 标记位 */ uint8 t_hoff; /* 包括NULL字段位图、对齐填充在内的元组头部大小 */ bits8 t_bits[FLEXIBLE_ARRAY_MEMBER]; /* NULL字段位图 */ /* 实际元组数据再该元组头部结构体之后,距离元组头部处偏移t_hoff字节 */ } HeapTupleHeaderData; 该结构体只是元组头部的定义,元组内容跟在该结构体后面,距离元组头部起始处的偏移由“t_hoff”成员保存。上面元组头部结构体部分成员信息,同时也构成了该元组的系统字段(字段序号小于0的那些字段)。对各个结构体成员的含义说明如下。 (1) t_xmin,插入元组的事务号(32位)。对应系统字段序号是MinTransactionIdAttributeNumber(-3)。 (2) t_xmax,删除元组的事务号(32位)。如果元组还没有被删除,那么为零。对应系统字段序号MaxTransactionIdAttributeNumber(-5)。 (3) t_cid,插入或删除元组的命令号。对应系统字段序号MinCommandIdAttributeNumber(-4)和MaxCommandIdAttributeNumber(-6)。 (4) t_ctid,当前元组的页面和页面内元组指针下标。如果该元组被更新,为更新后元组的页面号和页面内元组指针下标。 (5) t_infomask2,元组属性掩码,包含元组中字段个数、HOT(heap only tuple,堆内元组)更新标记、HOT元组标记等。 (6) t_infomask,元组另一个属性掩码,包含是否有空字段标记、是否有变长字段标记、是否有外部TOAST(the oversized-attribute storage technique,过长字段存储技术)标记、是否有OID字段标记、是否有压缩标记、插入事务是否提交/回滚标记、删除事务是否提交/回滚标记、是否被更新标记等。如果OID标记存在,那么元组OID从“t_hoff”偏移位置之前4个字节获得,对应系统字段序号ObjectIdAttributeNumber(-2)。 (7) t_hoff,元组数据距离元组头部结构体起始位置的偏移。 (8) t_bits,所有字段的NULL值bitmap。每个字段对应t_bits中的一个bit位,因此是变长数组。 上述元组结构体在内存中使用时嵌入在一个更大的元组数据结构体中,该结构体的定义代码如下。除了保存元组内容的t_data成员之外,其他的成员保存了该元组的一些其他系统信息,这些信息构成了该元组剩余的一些系统字段内容: typedef struct HeapTupleData { uint32 t_len; /* 包括元组头部和数据在内的元组总大小 */ ItemPointerData t_self; /* 元组行号 */ Oid t_tableOid; /* 元组所属表的OID */ TransactionId t_xid_base; TransactionId t_multi_base; HeapTupleHeader t_data; /* 指向元组头部 */ } HeapTupleData; 该结构体主要成员的含义如下。 (1) t_len,元组长度。 (2) t_self,元组所在页面号和页面内元组指针下标,对应系统字段序号SelfItemPointerAttributeNumber(-1)。 (3) t_tableOid,该元组所属表的OID,对应系统字段序号TableOidAttributeNumber(-7)。 介绍了astore堆表元组结构,下面介绍常用的astore堆表元组操作接口。如表4-11所示。 表4-11 常用的元组操作接口 操作接口名 操作含义 对应的行存储统一访存接口 heap_form_tuple 使用传入的、各个元组字段的values数组和nulls数组,生成一条完整的元组。一般用于插入操作 tableam_tops_form_tuple heap_deform_tuple 使用传入的完整元组以及各个字段的类型定义,解构各个字段的值,生成values数组和nulls数组。一般用于更新前的准备工作 tableam_tops_deform_tuple heap_modify_tuple 先调用heap_deform_tuple解构传入的原始元组,然后将解构得到的values和nulls数组中需要更新的字段替换为新的值,最后再调用heap_form_tuple生成修改后的完整元组。一般用于更新操作 tableam_tops_modify_tuple heap_freetuple 释放一条元组对应的内存空间 tableam_tops_free_tuple heap_copytuple 复制一条完整的元组,包括元组头和元组内容 tableam_tops_copy_tuple heap_form_cmprs_tuple 类似heap_form_tuple,生成一条压缩后的元组 tableam_tops_form_cmprs_tuple heap_deform_cmprs_tuple 类似heap_deform_tuple,解构一条压缩后的元组 tableam_tops_deform_cmprs_tuple heap_getattr 获取一条元组中指定的用户或系统字段值 tableam_tslot_getattr heap_getsysattr 获取一条元组中指定的系统字段值 tableam_tops_getsysattr 在上述操作接口中,heap_getattr操作接口是最常用的操作接口之一,执行流程如图4-3所示。 图4-3 heap_getattr操作接口从元组中获取单个字段值的流程图 heap_getattr操作接口在代码上做了多处优化: (1) 判断待访问的字段序号是否大于元组头部保存的元组实际字段个数;如果大于,则通过访问pg_attribute系统表得到。该优化来自快速追加表字段特性。该特性允许用户在不需要重写一张表所有行的情况下,在一张表的最后增加一个或多个带默认值约束的字段。 (2) 如果该元组的字段全部非空并且待查询字段之前所有的字段都是定长的,那么在上一个heap_getattr查询该字段的操作过程中,会缓存该字段在元组中的字节偏移;之后再次查询时,当满足元组字段全部非空的情况下会使用上述缓存的偏移位置直接读取字段内容。 (3) 读取元组头部的NULL值bitmap,如果该字段对应的bitmap中的比特位非0,则直接返回NULL值。 2) astore堆表页面结构 由于整体行存储格式默认的介质管理器是磁盘文件系统,因此采用了和文件系统类似的段页式设计,最小I/O单元为一个页面,这样可以在大多数场景下获得比较好的I/O性能和较低的I/O开销。一个astore堆表页面默认大小为8kB,其结构如图4-4所示。 图4-4 astore堆表页面结构示意图 在一个astore堆表页面中,页面头部分对应HeapPageHeaderData结构体。其中,pd_multi_base以及之前的部分对应定长成员,存储了整个页面的重要元信息;pd_multi_base之后的部分对应元组指针变长数组,其每个数组成员存储了页面中从后往前的、每个元组的起始偏移和元组长度。如图4-4所示,真正的元组内容从页面尾部开始插入,向页面头部扩展;相应的,记录每条元组的元组指针从页面头定长成员之后插入,往页面尾部扩展;整个页面中间形成一个空洞,供后续插入的元组和元组指针使用。 对于一个astore堆表的一条具体元组,有一个全局唯一的逻辑地址,即元组头部的t_ctid,其由元组所在的页面号和页面内元组指针数组下标组成;该逻辑地址对应的物理地址,则由ctid和对应的元组指针成员共同给出。通过页面、对应元组指针数组成员、页面内偏移和元组长度的访问顺序,就可以完整获取到一条元组的完整内容。t_ctid结构体和元组指针结构体的定义代码如下。 /* t_ctid结构体*/ typedef struct ItemPointerData { BlockIdData ip_blkid; /* 页号 */ OffsetNumber ip_posid; /* 页面偏移,即对应的页内元组指针下标 */ } ItemPointerData; /* 页面内元组指针结构体 */ typedef struct ItemIdData { unsigned lp_off : 15, /* 元组起始位置(距离页头) */ lp_flags : 2, /* 元组指针状态 */ lp_len : 15; /* 元组长度 */ } ItemIdData; 如上两级的元组访问设计,主要有两个优点。 (1) 在索引结构中(参见“4.2.5 行存储索引机制”小节),只需要保存元组的t_ctid值即可,无须精确到具体字节偏移,从而降低了索引元组的大小(节约两个字节),提升索引查找效率; (2) 将页面内元组的地址查找关系自封闭在页面内部的元组指针数组中,和外部索引解耦,从而在某些场景下可以让页面级空闲空间整理对外部索引数据没有影响,降低空闲空间回收的开销和设计复杂度。具体实现机制在“5. astore空间管理和回收”小节中介绍。 astore堆表页面头具体结构体定义代码如下: typedef struct { PageXLogRecPtr pd_lsn; /* 页面最新一次修改的日志lsn */ uint16 pd_checksum; /* 页面CRC */ uint16 pd_flags; /* 标志位 */ LocationIndex pd_lower; /* 空闲位置开始出(距离页头) */ LocationIndex pd_upper; /* 空闲位置结尾处(距离页头) */ LocationIndex pd_special; /* 特殊位置起始处(距离页头) */ uint16 pd_pagesize_version; ShortTransactionId pd_prune_xid; TransactionId pd_xid_base; TransactionId pd_multi_base; ItemIdData pd_linp[FLEXIBLE_ARRAY_MEMBER]; } HeapPageHeaderData; 其中各个成员的含义如下。 (1) pd_lsn:该页面最后一次修改操作的预写日志结束位置的下一个字节,用于检查点推进和保持恢复操作的幂等性(幂等指对接口的多次调用所产生的结果和调用一次是一致的)。 (2) pd_checksum:页面的CRC校验值。 (3) pd_flags:页面标记位,用于保存各类页面相关的辅助信息,如页面是否有空闲的元组指针、页面是否已满、页面元组是否都可见、页面是否被压缩、页面是否是批量导入的、页面是否加密、页面采用的CRC校验算法等。 (4) pd_lower:页面中间空洞的起始位置,即当前已使用的元组指针数组的尾部。 (5) pd_upper:页面中间空洞的结束位置,即下一个可以插入元组的起始位置。 (6) pd_special:页面尾部特殊区域的起始位置。该特殊位置位于第一条元组记录和页面结尾之间,用于存储一些变长的页面级元信息,如采用的压缩算法信息、索引的辅助信息等。 (7) pd_pagesize_version:页面的大小和版本号。 (8) pd_prune_xid:页面清理辅助事务号(32位),通常为该页面内现存最老的删除或更新操作的事务号,用于判断是否要触发页面级空闲空间整理。实际使用的64位prune事务号由“pd_prune_xid”字段和“pd_xid_base”字段相加得到。 (9) pd_xid_base:该页面内所有元组的基准事务号(64位)。该页面所有元组实际生效的64位xmin/xmax事务号由“pd_xid_base”(64位)和元组头部的“t_xmin/t_xmax”字段(32位)相加得到。 (10) pd_multi_base:类似“pd_xid_base”字段,当对元组加锁时,会将持锁的事务号写入元组中,该64位事务号由“pd_multi_base”字段(64位)和元组头部的“t_xmax”字段(32位)相加得到。 (11) pd_linp:元组指针变长数组。 对于astore堆表页面的主要管理接口如表4-12所示。鉴于astore采用的元组多版本设计实现方式(参见“3. astore元组多版本机制”小节),删除操作并不会直接从页面中删除指定的元组,页面管理也没有提供这样的接口。对于被删除的、过于陈旧的元组,通过页面空闲空间整理流程(参见“5. astore空间管理和回收”小节)完成。 表4-12 页面管理接口函数 函数名 操作含义 PageAddItem 在页面中插入一条新的元组 PageRepairFragmentation 页面空闲空间整理 在astore堆表页面中,采用64位页面“pd_xid_base”字段和32位元组“t_xmin/t_xmax”字段组合设计方式的原因如下。 早期openGauss版本采用32位事务号,所以对于OLTP类系统事务号消耗速度很快。当消耗的事务号超过最大事务号一半左右时,整个系统会强制对所有元组进行防止事务号回卷的整理工作。这个过程将阻塞所有写查询,系统不可用。 为了解决这个问题,openGauss将事务号升级到64位。为了平滑兼容之前32位事务号的元组头部结构,没有改变元组的结构和长度,而是在32位事务号页面头部结构体的基础上,扩展增加了标识整个页面所有元组事务号范围的64位基准事务号“pd_xid_base”和“pd_multi_base”两个字段。同一个页面中所有元组实际的64位“xmin/xmax”字段,一定在“pd_xid_base”字段和“pd_xid_base+2322”之间。 可以通过astore堆表页面头部“pd_pagesize_version”字段中页面版本号来区分32位事务号页面和64位事务号页面: (1) 版本号等于4,为32位事务号页面。 (2) 版本号等于5,为64位非堆表页面(如索引页面)。这类页面的页头无须保存64位事务号信息,因此和32位事务号页面采用相同的结构。这类页面中可能涉及的64位事务号信息,保存在页面尾部的“”pd_special”字段区域中。 (3) 版本号等于6,即为64位astore堆表页面。 对于从32位事务号系统升级上来的astore堆表页面,在部分页面访问场景中(如RelationGetBufferForTuple/heap_delete/heap_update/heap_lock_tuple),首先会判断访问的页面是否是4号版本。若是4号版本,则调用heap_page_upgrade尝试进行页面版本升级。当页面空闲空间足够放下扩展的两个成员(共16个字节)时,调用PageLocalUpgrade函数将页面格式升级到64位,且升级后的pd_xid_base字段和pd_multi_base字段一定为0;如果剩余空间不够,系统会给出报错或告警,并提示用户执行VACUUM FULL命令来手动升级页面。 对于需要修改元组事务号的操作(如heap_insert/heap_multi_insert/heap_delete/heap_update/heap_lock_tuple),需要判断新写入的64位事务号是否满足在页面的“pd_xid_base”和“pd_xid_base+232”之间。如果满足,则通过检查;否则,需要调整页面的“pd_xid_base”字段或“pd_multi_base”字段的值以满足上述条件。如果新写入的事务号和页面上现有任意一个元组的“xmin/xmax”事务号差距已经超过232,系统还会尝试对现有元组进行“freeze”(冻结)操作。如果“freeze”操作之后,上述事务号差距还是超过232,该查询会报错退出。 32位事务号astore堆表页面头结构代码如下所示,各成员含义可参考64位事务号页面头结构: typedef struct { PageXLogRecPtr pd_lsn; uint16 pd_checksum; uint16 pd_flags; LocationIndex pd_lower; LocationIndex pd_upper; LocationIndex pd_special; uint16 pd_pagesize_version; ShortTransactionId pd_prune_xid; ItemIdData pd_linp[FLEXIBLE_ARRAY_MEMBER]; } PageHeaderData; 下篇我们将详细介绍“3. astore元组多版本机制”相关内容,敬请期待!

优秀的个人博客,低调大师

openGauss数据库源码解析系列文章——存储引擎源码解析(一)

OLTP、OLAP业务各自对数据库的存储引擎提出了不同的要求,而openGauss能够支持多个存储引擎来满足来自不同场景的业务诉求。本章将逐一介绍各种存储引擎和对应的源码。 4.1 存储引擎整体架构及代码概览 从整个数据库服务的组成构架来看,存储引擎向上对接SQL引擎,为SQL引擎提供或接收标准化的数据格式(元组或向量数组);向下对接存储介质,按照特定的数据组织方式,以页面、列存储单元(CU,compression unit)或其他形式为单位,通过存储介质提供的特定接口,对存储介质中的数据完成读、写操作。在此基础之上,存储引擎通过日志系统提供数据的持久化和可靠性能力;通过并发控制(事务)系统保证同时执行的、多个读写操作之间的原子性、一致性和隔离性;通过索引系统提供对特定数据的加速寻址和查询能力;通过主备复制系统提供整个数据库服务的高可用能力。 图4-1 openGauss存储引擎整体构架示意图 图4-1是openGauss存储引擎整体构架的示意图。总体上,根据存储介质和并发控制机制,存储引擎分为磁盘引擎和内存引擎两大类。磁盘引擎主要面向通用的、大容量的业务场景,内存引擎主要面向容量可控的、追求极致性能的业务场景。在磁盘引擎中,为了满足不同业务场景对于数据不同的访问和使用模式,openGauss进一步提供了astore(append-store,追加写优化格式)、cstore(column store,列存储格式)以及可拓展的数据元组和数据页面组织格式。在内存引擎中,openGauss当前提供基于Masstree结构组织的mstore(memory-store,内存优化格式)数据组织格式。 上述几种引擎和存储格式的介绍如表4-1所示。 表4-1 openGauss存储引擎种类 父类 (存储介质和并发控制) 子类 (数据组织形式) 说明 磁盘引擎 (磁盘介质, 多版本和悲观并发控制(pessimistic concurrency control,PCC)) astore (追加写优化格式) 主要面向通用的在线交易处理类业务应用场景,适合高并发、小数据量的单点或小范围数据读、写操作。astore为行存储格式,向上提供元组形式的读、写;向下以页面为单位通过可扩展的介质管理器对存储介质进行读、写操作;并通过页面粒度的共享缓冲区来优化读、写操作的效率。当前行存存储格式默认的介质管理器采用磁盘文件系统接口,后续可扩展支持块设备等其他类型的存储介质 cstore (列存储格式) cstore (列存储格式) 面向联机分析处理类业务应用场景,适合大数据量的复杂查询和数据导入。cstore为列存存储格式,向上提供向量数组形式的读、写接口;向下以压缩单元为单位将数据保存在磁盘文件系统中(当前列存存储格式唯一支持的存储介质)。考虑到联机分析处理类业务通常以读操作为主,因此还提供了以压缩单元为粒度的只读共享缓冲区,以加速压缩单元的读操作性能 扩展存储格式 扩展存储格式 对于行存储类存储格式,openGauss提供了与上层SQL引擎对接的、统一的、可扩展的访存接口层(table access method)。该行存储统一访存接口层为SQL引擎提供元组形式的读、写接口,同时屏蔽了下层各种不同行存储类存储格式的内部实现,从而实现了SQL引擎与存储引擎(行存储类磁盘引擎)的解耦,大幅提升了不同存储格式之间的隔离性和开发效率 当前行存储类存储格式支持追加写优化的astore格式,后续会支持更新写优化的ustore格式等其他数据组织格式 内存引擎 (内存介质, 乐观并发控制(OCC,optimistic concurrency control)) mstore (内存存储格式) mstore内存引擎面向超低时延和超高吞吐量的OLTP场景。数据以元组粒度存储于内存介质中,得益于内存介质读、写操作的超低时延(与磁盘介质相比),内存引擎可以提供极致的OLTP业务性能。内存引擎通过openGauss的外表访存接口实现与SQL引擎的数据交互 有如下几个特点。 (1) 统一的日志系统。 在openGauss的存储引擎中,磁盘引擎和内存引擎共用同一套日志系统,以保证在数据库故障恢复场景下,各个引擎内和各个引擎间的数据持久性和一致性。基于上述统一的日志系统,openGauss支持主、备机(主、备数据库服务进程)之间的流式日志复制,并通过Quorum复制协议,在保证复制一致性的前提下,尽可能降低日志同步对主机业务的影响。 (2)多种并发控制和事务系统。 在openGauss的存储引擎中,有两种并发控制和事务系统:适合高并发、高冲突、追求确定性结果的悲观并发控制机制;适合低冲突、短平快、低时延的乐观并发控制机制。 在磁盘引擎中,采用读写冲突优化的悲观并发控制机制:对于读、写并发操作,采用多版本并发控制(MVCC,multi-version concurrency control);对于写、写并发操作,采用基于两阶段锁协议(2PL,two-phase locking)的悲观并发控制(PCC,pessimistic concurrency control)。 在内存引擎中,采用乐观并发控制来尽可能降低并发控制系统对业务的阻塞,以获得极致的事务处理性能和时延。 (3) 表级存储格式/存储引擎和跨格式事务。 在openGauss的存储引擎中,支持在建表语句中指定目标表的存储格式和存储引擎,即行存储astore、列存储cstore、内存mstore和后续扩展的其他存储格式或存储引擎。因此,在同一个数据库中,为了适配不同的业务场景,用户可以创建不同存储格式或不同存储引擎的表。进一步,当前openGauss在同一个事务内,支持对同一引擎不同存储格式的表的读写查询,这将极大地简化不同存储格式表中数据一致性、同步性和实时性的运维难度。后续openGauss版本计划支持跨引擎事务,这将使得openGauss数据库在面对多样化的业务场景时显得更为游刃有余。 (4)统一的行存储访存接口。 在openGauss的磁盘引擎中,行存储类存储格式是最传统也是使用场景最广泛的存储格式。针对不同的业务场景,行存储格式需要进行不同的优化和设计。为了便于后续新型行存储格式的扩展,在openGauss中提供了统一的行存储访存接口层,为上层SQL引擎屏蔽了底层不同的行存储数据组织形式。 对于不同的行存储数据格式,它们向上对接统一的行存储访存接口,向下共享缓冲区管理、事务并发控制、日志系统、持久化和故障恢复、主备系统、索引机制。同时,不同的行存储数据格式内部又实现了不同的元组和页面格式,以及在此之上的访存接口、元组多版本、页面多版本、空闲空间管理回收等不同功能。 openGauss存储引擎的代码主要位于“src/gausskernel/storage/”目录下,具体目录结构如下: --src --gausskernel --storage --access --buffer --bulkload --cmgr --cstore --dfs --file --freespace --ipc --large_object --lmgr --mot --page --remote --replication --smgr 每个子目录都是一个相对独立的模块,和本章内容相关的如表4-2所示。 表4-2 存储引擎子目录 模块名 子目录 说明 访存模块 access子目录 主要包括:各种行存储格式中,元组格式;元组与页面之间的转换和访存管理;元组扫描、插入、删除和更新功能的接口实现;几类索引,包括B-Tree、hash、GIN(generalized inverted index,通用倒排索引)、GiST(generalized search tree,通用搜索树)、psort(列存储局部排序索引),的访存管理和接口实现;各类数据库操作对应的日志实现和恢复机制;以及事务模块实现 行存储共享缓冲区模块 buffer子目录 主要包括:行存储共享缓冲区的结构;物理页面和缓冲区页面的映射管理;缓存页面的加载和淘汰算法等 列存储只读共享缓冲区模块 cmgr子目录 主要包括:cstore列存储格式只读共享缓冲区的结构;压缩单元和缓冲区的映射管理;缓冲压缩单元的加载和淘汰算法等 列存储访存模块 cstore子目录 主要包含:cstore列存储格式中,向量数组与压缩单元之间的转换和访存管理;以及在此基础之上向量数组的扫描、插入、删除和更新功能的接口实现 文件操作和虚拟文件描述符模块 file子目录 主要包含:磁盘文件系统存储介质的文件和目录操作;虚拟文件描述符的实现和管理 行存储空闲空间管理模块 freespace子目录 主要包含:各种行存储格式中,页面空闲空间的管理 内存引擎模块 mot子目录 主要包含:内存引擎的实现 页面模块 page子目录 主要包含:各种行存储格式中,页面格式、页面校验、页面加密和页面压缩 备机页面修复模块 remote子目录 主要包含:从备机获取完整页面或压缩单元,用于修复主机损坏的页面或压缩单元 主备日志复制模块 replication子目录 主要包含:主备日志发送和接收线程的实现;流式日志同步功能的实现;Quorum复制协议的实现,逻辑日志的实现以及主备重建;主备心跳检测功能的实现 存储介质管理模块 smgr子目录 主要包含:存储介质管理层的实现;磁盘文件系统(当前默认的存储介质)的基本功能接口实现 除了以上的这些模块之外,storage目录下剩余的子目录分别属于:外表批量导入模块(bulkload子目录)、外表服务器连接模块(dfs子目录)、进程间通信模块(ipc子目录)、大对象模块(large\_object子目录)、锁管理模块(lmgr子目录)。 openGauss存储引擎相关的后台线程实现代码包含在“src/gausskernel/process/postmaster”目录下,简要介绍如表4-3。在后序介绍具体相关模块消息序列时会详细介绍这些线程的工作原理和执行流程。 表4-3 存储引擎后台线程 线程名 文件名 说明 ADIO线程 aiocompleter.cpp 该线程主要负责异步-同步读写操作(ADIO,asynchronous-direct input-ouput)的后台预取和回写 autovacuum线程 autovacuum.cpp 该线程主要负责磁盘引擎的后台空闲空间回收 bgwriter线程 bgwriter.cpp 该线程主要负责行存储表的后台脏页写入磁盘(当内存数据页跟磁盘数据页内容不一致的时候,称这个内存页为“脏页”。内存数据写入到磁盘后,内存和磁盘上的数据页的内容就一致了,称为“干净页”) cbmwriter线程 cbmwriter.cpp 该线程主要负责增量页面修改信息的后台异步提取和CBM(changed block map,修改页面位图)日志的记录 checkpointer线程 checkpointer.cpp 该线程主要负责在后台定期推进数据库的故障恢复点 lwlockmonitor线程 lwlockmonitor.cpp 该线程主要负责业务线程轻量级锁的死锁检测 pagewriter线程 pagewriter.cpp 该线程主要负责行存储共享缓冲区的脏页写入磁盘 pgarch线程 pgarch.cpp 该线程主要负责在后台定期执行日志归档命令 remoteservice线程 remoteservice.cpp 该线程主要负责接收主机页面修复RPC(remote procedure call,远程函数调用)请求 startup线程 startup.cpp 该线程为数据库故障恢复和回放日志的主线程 walwriter线程 walwriter.cpp 该线程主要负责在后台异步写入磁盘日志 4.2 磁盘引擎 磁盘引擎是数据库系统中最常用的存储引擎,openGauss提供不同存储格式的磁盘引擎来支持大容量(数据量大于内存空间)场景下的OLTP、OLAP和HTAP(hybrid transactions and analytics processing,混合交易和分析处理)业务。本节主要介绍openGauss数据库内核中磁盘引擎的实现方式。 4.2.1 磁盘引擎整体框架及代码概览 磁盘引擎的整体框架如图4-1中所示。根据与上层SQL引擎之间交互的数据结构类型,可以分为行存储格式和列存储格式。这两种数据格式共用相同的事务并发控制、日志系统、持久化和故障恢复、主备系统。 在此基础之上,行存储格式内部设计为可以支持多种不同子格式的可扩展架构。不同行存储子格式之间共用相同的行存储统一访存接口(table access method)、共享缓冲区、索引机制等。当前仅支持追加写优化的astore子格式,后续计划支持写优化的ustore子格式以及面向其他场景优化的其他子格式。另一方面,在openGauss行存储格式中,对同一行数据的写-写查询冲突通过两阶段锁协议来实现并发控制(参见第5章中关于行级锁的介绍),对同一行数据的读-写查询冲突通过行级多版本技术来实现互不阻塞的、高效的并发控制。对于不同的行存储子格式,可能采用不同的行级多版本实现方式,从而也会引入不同的、清理历史版本的空闲空间管理和回收机制。 磁盘引擎的主要功能模块和代码分布如表4-4所示。 表4-4 磁盘引擎功能模块 功能模块名 说明 行存储统一访存管理 向上对接SQL引擎,提供对行存储表各类访存操作的抽象接口,包括:行级查询、插入、删除、修改等操作接口;向下根据行存储表实际的行存储子格式,调用与子格式对应的具体访存操作实现 代码主要在“src/gausskernel/storage/access/table”目录下 astore访存管理 提供astore行存储格式表的具体访存操作实现,包括:对astore堆表的行级查询、插入、删除、修改等操作接口;astore堆表行级多版本机制和元组可见性判断;根据astore堆表页间、页内结构,以及astore堆表元组结构,完成对astore堆表文件的遍历和增删改查操作 代码主要在“src/gausskernel/storage/access/heap”目录(单表文件管理)和“src/gausskernel/storage/access/hbstore”目录(段页式文件管理)下 astore堆表/索引表页面结构 包括astore堆表/索引表元组在页面内的具体组织形式,在页面内插入元组操作、页面整理操作、页面初始化、页面加解密、页面CRC(cyclic redundancy check,循环冗余码校验)校验操作等 代码主要在“src/gausskernel/storage/access/redo/bufpage.cpp”文件、“redo_bufpage.cpp”文件和对应头文件中 astore堆表元组结构 包括astore堆表元组的结构、填充、解构、修改、字段查询、变形、压缩、解压等操作 代码主要在“src/gausskernel/storage/access/common/heaptuple.cpp”文件和对应头文件中 行存储索引访存管理 向上对接SQL引擎,提供对索引表的行级查询、插入、删除等操作接口;向下根据索引表页间、页内结构,以及索引表元组结构,完成对指定索引键的查找和增删操作 索引访存层抽象框架代码在“src/gausskernel/storage/access/index”目录下,每种索引结构具体对应的实现代码在同级的gin目录、gist目录、hash目录、nbtree目录、spgist目录 行存储索引表元组结构 包括行存储索引表元组的结构、填充、解构、拷贝等操作 代码主要在“src/gausskernel/storage/access/common/indextuple.cpp”文件和对应头文件中 行存储共享缓冲区管理 包括共享缓冲区的结构、页面查找方式、页面淘汰方式等 代码主要在“src/gausskernel/storage/buffer”目录下 行存储介质管理器管理和堆表/索引表文件管理 包括几种主要介质操作的抽象接口以及几种主要的、基于磁盘文件系统的堆表/索引表文件操作接口 代码在“src/gausskernel/storage/smgr”目录下 cstore访存管理 向上对接SQL引擎,提供对cstore列存储表的向量数组(vector batch)粒度的查询、插入、删除、修改等操作接口;向下根据cstore列存储表CU间、CU内结构,完成对cstore列存储表文件的遍历和增删改查操作;cstore列存储表CU内和CU间的多版本并发控制和可见性判断 代码主要在“src/gausskernel/storage/cstore”目录下的cstore_系列文件中 cstore索引访存管理 向上对接SQL引擎,提供对cstore索引表的向量数组粒度的查询、插入等操作接口;向下根据cstore索引表组织结构,完成对指定索引键的查询和插入等操作 代码主要在“src/gausskernel/storage/access/cbtree”目录(cstore列储存B-Tree索引)下和“src/gausskernel/storage/access/psort”目录(cstore列存储psort索引)下 cstore列存储表 CU结构 ① 和行存不同,cstore列存储表与外存的I/O单元为CU。该部分主要包括CU的内部结构、CU的填充和压缩等操作 ② 代码在“src/gausskernel/storage/cstore/cu.cpp”文件中 cstore列存储表 CU只读共享缓冲区管理 包括以CU为单位的只读共享缓冲区的结构、查找、淘汰等 代码主要在“src/gausskernel/storage/cmgr”目录下 cstore列存储表 CU持久化介质模块 包括以CU为粒度的、基于磁盘介质的cstore列存储表文件外存I/O操作 代码在“src/gausskernel/storage/cstore/custorage.cpp”文件中 预写日志共享缓冲区和文件管理 包括日志记录格式、日志页面格式、日志文件格式、日志插入、日志写入磁盘、日志缓冲区管理、日志归档、日志恢复等操作 代码在“src/gausskernel/storage/access/transam/xlog”系列文件中 检查点和故障恢复管理 包括页面淘汰算法和检查点推进算法、双写刷盘(写入磁盘)、页面故障恢复等 代码主要分布在“src/gausskernel/process/postmaster/pagewriter.cpp”、“src/gausskernel/process/postmaster/bgwriter.cpp”、“src/gausskernel/storage/access/transam/double_write.cpp”、“src/gausskernel/storage/access/transam/xlog.cpp”、对应头文件和“src/gausskernel/storage/access/redo”目录 事务管理和并发控制 包括锁管理、事务提交流程、快照维护、提交时间戳维护、可见性判断等 代码主要在“src/gausskernel/storage/access/transam”目录下 该部分内容较为复杂,在第5章单独介绍 事务提交日志SLRU(Simple Least Recently Used,简单最近最少使用)共享缓冲区和文件管理 包括事务提交日志的页面格式、读写操作、SLRU缓存算法、清理操作等,与事务管理模块一起介绍 事务提交时间戳日志SLRU共享缓冲区和文件管理 包括事务提交日志(CSNLOG)的页面格式、读写操作、SLRU缓存算法、清理操作等,与事务管理模块一起介绍 关键控制文件管理 主要包括控制文件、根系统表文件等关键文件的读、写操作 代码分布较广 在上述模块基础之上,openGauss磁盘引擎还包括CU压缩、外表、批量导入等功能,代码分布在“src/gausskernel/storage/cstore/compression”、“src/gausskernel/storage/access/dfs”、“src/gausskernel/storage/bulkload”等目录下。 openGauss磁盘引擎的关键技术整体来说包括: (1) 基于事务提交逻辑时间戳的快照隔离机制以及多版本并发控制技术。 (2) 基于事务号(xid,全称transaction identifier)的行级多版本管理技术。 (3) 基于大内存设计的共享缓冲区管理和淘汰算法。 (4) 平滑无性能波动的增量检查点(checkpoint)技术。 (5) 基于并行回放的快速故障实例恢复技术。 (6) 支持事务语义的DML操作和DDL操作。 (7) 面向OLAP场景的cstore列存储格式。当表中列数比较多、但是访问的列数比较少时可以大大减少不必要的列的I/O开销。 (8) 面向OLAP场景的cstore列存储批量访存接口。向上支持以向量数组为粒度的批量数据访存接口,结合向量化执行引擎提升CPU缓存命中率和系统吞吐率。 (9) 面向OLAP场景的cstore列存储高效压缩算法。基于同一列比较相似的数据特征,在大数据量下获得很高的压缩效果,减少系统的I/O开销。 4.2.2 行存储统一访存接口 如上所述,在openGauss中,提供行存储统一访存接口层,来屏蔽不同行存储子格式内部实现机制对SQL引擎的影响。该行存储统一访存接口层被称为Table Access Method层。根据SQL引擎对行存储表的访存方式,将访存接口分为5类,如表4-5所示。每一类接口的具体操作如表4-6至4-10所示。 表4-5 Table Access Method定义的访存接口 接口类别 接口含义 Tuple AM Slot AM 元组(tuple)和元组槽(slot)操作抽象层,包括元组数据结构的抽象、元组操作的抽象,执行引擎无须关注元组属于哪种行存储子格式,只需调用元组数据结构基类的抽象操作接口,就可操作不同行存储子格式的元组,从而屏蔽不同行存储子格式物理元组结构、访问方法的差异 TableScan AM 表扫描(table scan)抽象层,包括TableScan数据结构的抽象、TableScan管理操作的抽象,执行引擎无须关注行存储子格式内部TableScan结构的差异,通过调用TableScan数据结构基类的抽象管理接口,就可完成不同行存储子格式的TableScan管理,屏蔽不同行存储子格式内部实现的差异 DQL AM 元组查询(data query language,DQL)操作抽象层,包括获取元组、元组可见性判断等查询操作的抽象 DML AM 元组写操作抽象层,包括元组插入、批插、删除、更新、锁定等接口的抽象 DDL AM 表物理操作抽象层,这里统称为DDL抽象层,涉及表物理文件操作的相关接口的抽象,例如CTAS、TRUNCATE、LOAD/COPY、VACUUM、VACUUM FULL、ANALYZE、REBUILD INDEX、ALTER TABLE RESTRUCT等DDL语法。该层也可以支持存储管理的抽象功能,如屏蔽不同行存储子格式的文件/目录管理模块、SMGR访问等差异 表4-6 Tuple AM、Slot AM类访存接口 接口名称 接口含义 tableam_tslot_clear 清理slot tuple,主要是被ExecClearTuple调用 tableam_tslot_materialize 该方法在ExecMaterializeSlot被调用, 将slot中的tuple进行local copy(本地拷贝) tableam_tslot_get_minimal_tuple 获取slot中的minimal tuple(最小化元组),slot负责管理/释放minimal tuple的内存 tableam_tslot_copy_minimal_tuple 返回slot中minimal tuple的副本,该副本在当前内存上下文中被分配,需要调用者进行释放操作 tableam_tslot_store_minimal_tuple 此函数在指定的TupleTableSlot结构体中存储minimal tuple tableam_tslot_get_heap_tuple 该函数获取slot中的tuple tableam_tslot_copy_heap_tuple 该函数返回slot中tuple的副本,该副本在当前内存上下文中被分配,需要调用者进行释放操作 tableam_tslot_store_tuple 该方法将对应的物理元组存储到slot中 tableam_tslot_getsomeattrs 强制更新slot中tuple某个属性的values和isnull数组信息 tableam_tslot_getattr 获取当前slot中tuple的某个属性信息 tableam_tslot_getallattrs 强制更新slot中tuple的values和isnull数组 tableam_tslot_attisnull 检查slot中tuple的属性是否为null tableam_tslot_get_tuple_from_slot 从slot中获取一个tuple,并根据relation结构体中行存储子格式信息转换为对应子格式的tuple tableam_tops_getsysattr 获取tuple的系统属性 tableam_tops_form_minimal_tuple 根据values和isnull数组内容,新建一个tuple tableam_tops_form_tuple 根据values和isnull数组内容,新建一个minimal tuple tableam_tops_form_cmprs_tuple 根据values和isnull数组内容,新建一个被压缩的tuple tableam_tops_deform_tuple 抽取指定tuple中的data数据到values和isnull数组 tableam_tops_deform_cmprs_tuple 抽取被压缩的tuple中的data数据到values和isnull数组 tableam_tops_computedatasize_tuple 计算需要构造的tuple的data区域的大小 tableam_tops_fill_tuple 根据values和isnull数组中的数据填充到tuple的data区域 tableam_tops_modify_tuple 根据一个旧tuple新建一个tuple并更新其values tableam_tops_free_tuple 释放一个tuple的内存 tableam_tops_tuple_getattr 获取tuple的某个属性信息 tableam_tops_tuple_attisnull 检查tuple的属性是否为null tableam_tops_copy_tuple 拷贝并返回一个tuple tableam_tops_copy_minimal_tuple 拷贝并返回一个minimal tuple tableam_tops_free_minimal_tuple 释放minimal tuple的内存 tableam_tops_new_tuple 新建一个tuple tableam_tops_destroy_tuple 销毁一个tuple tableam_tops_get_t_self 获取tuple中的self指针,指向自己在表中的位置 tableam_tops_exec_delete_index_tuples 删除索引的tuple tableam_tops_exec_update_index_tuples 更新索引的tuple tableam_tops_get_tuple_type 获取tuple属于哪种存储引擎 tableam_tops_copy_from_insert_batch copy from场景进行批量INSERT(插入) tableam_tops_update_tuple_with_oid 根据table OID(表的唯一标识号)更新tuple 表4-7 TableScan AM类访存接口 接口名称 接口含义 tableam_scan_begin 初始化scan结构体,准备执行table scan(全表扫描)算子 tableam_scan_begin_bm 准备执行bitmap scan(位图扫描)算子 tableam_scan_begin_sampling 初始化堆表(顺序)扫描操作 tableam_scan_getnexttuple 返回scan中的下一个tuple tableam_scan_getpage 获取scan中的下一页 tableam_scan_end 结束scan,并释放内存 tableam_scan_rescan 重置scan tableam_scan_restrpos 重置扫描位置 tableam_scan_markpos 记录当前扫描位置 tableam_scan_init_parallel_seqscan 初始化并行sequence scan(顺序扫描) 表4-8 DQL AM类访存接口 接口名称 接口含义 tableam_tuple_fetch 根据tid(元组物理位置)获取tuple tableam_tuple_satisfies_snapshot 指定元组对于快照是否可见 tableam_tuple_get_latest_tid 获取tid指向的当前snapshot(快照)可见的最新物理元组 表4-9 DML AM类访存接口 接口名称 接口含义 tableam_tuple_insert 插入一条元组到表中 tableam_tuple_multi_insert 插入多条元组到表中 tableam_tuple_delete 删除一条元组,返回并发冲突状态,由调用者根据并发冲突状态决定下步操作 tableam_tuple_update 更新一条记录,返回并发冲突状态,由调用者根据并发冲突状态决定下步操作 tableam_tuple_lock 锁定一条元组 tableam_tuple_lock_updated 解锁一条元组 tableam_tuple_check_visible 检查元组的可见性 tableam_tuple_abort_speculative 终止upsert操作的尝试插入操作,转为更新操作 表4-10 DDL AM类访存接口 接口名称 接口含义 tableam_index_build_scan 该方法用于创建索引的首次全表扫描 tableam_index_validate_scan 该方法用于并发创建索引的第二次全表扫描 tableam_relation_copy_for_cluster 将源表数据根据指定的聚簇方式复制到新表中 对于每一个行存储子格式,需要提供上述这五类访存接口的各自实现方式,并注册到g_tableam_routines全局行存储访存接口数组中。SQL引擎在调用某个访存接口时,根据Relation结构体中表的子格式类型(rd_tam_type成员),来调用对应的子格式访存接口。 由于内容较多,下篇我们将详细介绍“4.2.3 astore”相关内容,敬请期待!

优秀的个人博客,低调大师

React Hooks源码深度解析

作者:京东零售 郑炳懿 前言 React Hooks是React16.8 引入的一个新特性,它允许函数组件中使用state和其他 React 特性,而不必使用类组件。Hooks是一个非常重要的概念,因为它们提供了更简单、更易于理解的React开发体验。 React Hooks的核心源码主要包括两个部分:React内部的Hook管理器和一系列预置的Hook函数。 首先,让我们看一下React内部的Hook管理器。这个管理器是React内部的一个重要机制,它负责管理组件中的所有Hook,并确保它们在组件渲染期间以正确的顺序调用。 内部Hook管理器 示例: const Hook = { queue: [], current: null, }; function useState(initialState) { const state = Hook.current[Hook.queue.length]; if (!state) { Hook.queue.push({ state: typeof initialState === 'function' ? initialState() : initialState, setState(value) { this.state = value; render(); }, }); } return [state.state, state.setState.bind(state)]; } function useHook(callback) { Hook.current = { __proto__: Hook.current, }; try { callback(); } finally { Hook.current = Hook.current.__proto__; } } function render() { useHook(() => { const [count, setCount] = useState(0); console.log('count:', count); setTimeout(() => { setCount(count + 1); }, 1000); }); } render(); 在这个示例中,Hook对象有两个重要属性:queue和current。queue存储组件中所有Hook的状态和更新函数,current存储当前正在渲染的组件的Hook链表。useState和useHook函数则分别负责创建新的Hook状态和在组件中使用Hook。 预置 Hook 函数 useState Hook 以下是useState Hook的实现示例: function useState(initialState) { const hook = updateWorkInProgressHook(); if (!hook.memoizedState) { hook.memoizedState = [ typeof initialState === 'function' ? initialState() : initialState, action => { hook.queue.pending = true; hook.queue.dispatch = action; scheduleWork(); }, ]; } return hook.memoizedState; } 上述代码实现了useState Hook,其主要作用是返回一个state和更新函数的数组,state 初始值为initialState。 在这个实现中,updateWorkInProgressHook()函数用来获取当前正在执行的函数组件的 fiber 对象并判断是否存在对应的hook。它的实现如下: function updateWorkInProgressHook() { const fiber = getWorkInProgressFiber(); let hook = fiber.memoizedState; if (hook) { fiber.memoizedState = hook.next; hook.next = null; } else { hook = { memoizedState: null, queue: { pending: null, dispatch: null, last: null, }, next: null, }; } workInProgressHook = hook; return hook; } getWorkInProgressFiber()函数用来获取当前正在执行的函数组件的fiber对象,workInProgressHook则用来存储当前正在执行的hook对象。在函数组件中,每一个useState调用都会创建一个新的 hook 对象,并将其添加到fiber对象的hooks链表中。这个hooks链表是通过fiber对象的memoizedState属性来维护的。 我们还需要注意到在useState Hook的实现中,每一个hook对象都包含了一个queue对象,用来存储待更新的状态以及更新函数。scheduleWork()函数则用来通知React调度器有任务需要执行。 在React的源码中,useState函数实际上是一个叫做useStateImpl的内部函数。 下面是useStateImpl的源码: function useStateImpl<S>(initialState: (() => S) | S): [S, Dispatch<SetStateAction<S>>] { const dispatcher = resolveDispatcher(); return dispatcher.useState(initialState); } 可以看到,useStateImpl函数的作用就是获取当前的dispatcher并调用它的useState方法,返回一个数组,第一个元素是状态的值,第二个元素是一个dispatch函数,用来更新状态。这里的resolveDispatcher函数用来获取当前的dispatcher,其实现如下: function resolveDispatcher(): Dispatcher { const dispatcher = currentlyRenderingFiber?.dispatcher; if (dispatcher === undefined) { throw new Error('Hooks can only be called inside the body of a function component. (https://fb.me/react-invalid-hook-call)'); } return dispatcher; } resolveDispatcher函数首先尝试获取当前正在渲染的fiber对象的dispatcher属性,如果获取不到则说 明当前不在组件的渲染过程中,就会抛出一个错误。 最后,我们来看一下useState方法在具体的dispatcher实现中是如何实现的。我们以useReducer的 dispatcher为例,它的实现如下: export function useReducer<S, A>( reducer: (prevState: S, action: A) => S, initialState: S, initialAction?: A, ): [S, Dispatch<A>] { const [dispatch, currentState] = updateReducer<S, A>( reducer, // $FlowFixMe: Flow doesn't like mixed types [initialState, initialAction], // $FlowFixMe: Flow doesn't like mixed types reducer === basicStateReducer ? basicStateReducer : updateStateReducer, ); return [currentState, dispatch]; } 可以看到,useReducer方法实际上是调用了一个叫做updateReducer的函数,返回了一个包含当前状态和dispatch函数的数组。updateReducer的实现比较复杂,涉及到了很多细节,这里不再展开介绍。 useEffect Hook useEffect是React中常用的一个Hook函数,用于在组件中执行副作用操作,例如访问远程数据、添加/移除事件监听器、手动操作DOM等等。useEffect的核心功能是在组件的渲染过程结束之后异步执行回调函数,它的实现方式涉及到 React 中的异步渲染机制。 以下是useEffect Hook的实现示例: function useEffect(callback, dependencies) { // 通过调用 useLayoutEffect 或者 useEffect 方法来获取当前的渲染批次 const batch = useContext(BatchContext); // 根据当前的渲染批次判断是否需要执行回调函数 if (shouldFireEffect(batch, dependencies)) { callback(); } // 在组件被卸载时清除当前 effect 的状态信息 return () => clearEffect(batch); } 在这个示例中,useEffect接收两个参数:回调函数和依赖项数组。当依赖项数组中的任何一个值发生变化时, React会在下一次渲染时重新执行useEffect中传入的回调函数。 useEffect函数的实现方式主要依赖于React中的异步渲染机制。当一个组件需要重新渲染时,React会将所有的state更新操作加入到一个队列中,在当前渲染批次结束之后再异步执行这些更新操作,从而避免在同一个渲染批次中连续执行多次更新操作。 在useEffect函数中,我们通过调用useContext(BatchContext)方法来获取当前的渲染批次,并根据shouldFireEffect方法判断是否需要执行回调函数。在回调函数执行完毕后,我们需要通过clearEffect方法来清除当前effect的状态信息,避免对后续的渲染批次产生影响。 总结 总的来说,React Hooks的实现原理并不复杂,它主要依赖于React内部的fiber数据结构和调度系统,通过这些机制来实现对组件状态的管理和更新。Hooks能够让我们在函数组件中使用状态和其他React特性,使得函数组件的功能可以和类组件媲美。 除了useState、useEffect等hook,React还有useContext等常用的Hook。它们的实现原理也基本相似,都是利用fiber架构来实现状态管理和生命周期钩子等功能。 以上是hook简单实现示例,它们并不是React中实际使用的代码,但是可以帮助我们更好地理解hook的核心实现方式。

优秀的个人博客,低调大师

Function源码解析与实践

作者:陈昌浩 1 导读 if…else…在代码中经常使用,听说可以通过 Java 8 的 Function 接口来消灭 if…else…!Function 接口是什么?如果通过 Function 接口接口消灭 if…else…呢?让我们一起来探索一下吧。 2 Function 接口 Function 接口就是一个有且仅有一个抽象方法,但是可以有多个非抽象方法的接口,Function 接口可以被隐式转换为 lambda 表达式。可以通过 FunctionalInterface 注解来校验 Function 接口的正确性。Java 8 允许在接口中加入具体方法。接口中的具体方法有两种,default 方法和 static 方法。 @FunctionalInterfaceinterface TestFunctionService{ void addHttp(String url);} 那么就可以使用 Lambda 表达式来表示该接口的一个实现。 TestFunctionService testFunctionService = url -> System.out.println("http:" + url); 2.1 FunctionalInterface 2.1.1 源码 @Documented@Retention(RetentionPolicy.RUNTIME)@Target(ElementType.TYPE)public @interface FunctionalInterface {} 2.1.2 说明 上图是 FunctionalInterface 的注解说明。通过上面的注解说明,可以知道 FunctionalInterface 是一个注解,用来说明一个接口是函数式接口。 函数式接口只有一个抽象方法。 可以有默认方法,因为默认方法有一个实现,所以不是抽象的。函数接口的实例可以用 lambda 表达式、方法引用或构造函数引用创建。 FunctionalInterface 会校验接口是否满足函数式接口: 类型必须是接口类型,不能是注释类型、枚举或类。 只能有一个抽象方法。 可以有多个默认方法和静态方法。 可以显示覆盖 java.lang.Object 中的抽象方法。 编译器会将满足函数式接口定义的任何接口视为函数式接口,而不管该接口声明中是否使用 FunctionalInterface 注解。 3 Function 接口主要分类 Function 接口主要分类: Function:Function 函数的表现形式为接收一个参数,并返回一个值。 Supplier:Supplier 的表现形式为不接受参数、只返回数据。 Consumer:Consumer 接收一个参数,没有返回值。 Runnable:Runnable 的表现形式为即没有参数也没有返回值。 3.1 Function Function 函数的表现形式为接收一个参数,并返回一个值。 3.1.1 源码 @FunctionalInterfacepublic interface Function<T, R> { R apply(T t); default <V> Function<V, R> compose(Function<? super V, ? extends T> before) { Objects.requireNonNull(before); return (V v) -> apply(before.apply(v)); } default <V> Function<T, V> andThen(Function<? super R, ? extends V> after) { Objects.requireNonNull(after); return (T t) -> after.apply(apply(t)); } static <T> Function<T, T> identity() { return t -> t; }} 3.1.2 方法说明 apply:抽象方法。将此函数应用于给定的参数。参数 t 通过具体的实现返回 R。 compose:default 方法。返回一个复合函数,首先执行 fefore 函数应用于输入,然后将该函数应用于结果。如果任意一个函数的求值引发异常,则将其传递给组合函数的调用者。 andThen:default 方法。返回一个复合函数,该复合函数首先对其应用此函数它的输入,然后对结果应用 after 函数。如果任意一个函数的求值引发异常,则将其传递给组合函数的调用者。 identity:static 方法。返回一个始终返回其输入参数的函数。 3.1.3 方法举例 1)apply 测试代码: public String upString(String str){ Function<String, String> function1 = s -> s.toUpperCase(); return function1.apply(str);} public static void main(String[] args) { System.out.println(upString("hello!")); } 通过 apply 调用具体的实现。执行结果: 2)compose 测试代码: public static void main(String[] args) { Function<String, String> function1 = s -> s.toUpperCase(); Function<String, String> function2 = s -> "my name is "+s; String result = function1.compose(function2).apply("zhangSan"); System.out.println(result);} 执行结果 如结果所示:compose 先执行 function2 后执行 function1。 3)andThen 测试代码: public static void main(String[] args) { Function<String, String> function1 = s -> s.toUpperCase(); Function<String, String> function2 = s -> "my name is "+s; String result = function1.andThen(function2).apply("zhangSan"); System.out.println(result);} 执行结果: 如结果所示: andThen 先执行 function1 后执行 function2。 identity 测试代码: public static void main(String[] args) { Stream<String> stream = Stream.of("order", "good", "lab", "warehouse"); Map<String, Integer> map = stream.collect(Collectors.toMap(Function.identity(), String::length)); System.out.println(map);} 执行结果: 3.2 Supplier Supplier 的表现形式为不接受参数、只返回数据。 3.2.1 源码 @FunctionalInterfacepublic interface Supplier<T> { /** * Gets a result. * * @return a result */ T get();} 3.2.2 方法说明 get:抽象方法。通过实现返回 T。 3.2.3 方法举例 public class SupplierTest { SupplierTest(){ System.out.println(Math.random()); System.out.println(this.toString()); }} public static void main(String[] args) { Supplier<SupplierTest> sup = SupplierTest::new; System.out.println("调用一次"); sup.get(); System.out.println("调用二次"); sup.get();} 执行结果: 如结果所示:Supplier 建立时并没有创建新类,每次调用 get 返回的值不是同一个。 3.3 Consumer Consumer 接收一个参数,没有返回值。 3.3.1 源码 @FunctionalInterfacepublic interface Consumer<T> { void accept(T t); default Consumer<T> andThen(Consumer<? super T> after) { Objects.requireNonNull(after); return (T t) -> { accept(t); after.accept(t); }; }} 3.3.2 方法说明 accept:对给定参数 T 执行一些操作。 andThen:按顺序执行 Consumer -> after ,如果执行操作引发异常,该异常被传递给调用者。 3.3.3 方法举例 public static void main(String[] args) { Consumer<String> consumer = s -> System.out.println("consumer_"+s); Consumer<String> after = s -> System.out.println("after_"+s); consumer.accept("isReady"); System.out.println("========================"); consumer.andThen(after).accept("is coming");} 执行结果: 如结果所示:对同一个参数 T,通过 andThen 方法,先执行 consumer,再执行 fater。 3.4 Runnable Runnable:Runnable 的表现形式为即没有参数也没有返回值。 3.4.1 源码 @FunctionalInterfacepublic interface Runnable { public abstract void run();} 3.4.2 方法说明 run:抽象方法。run 方法实现具体的内容,需要将 Runnale 放入到 Thread 中,通过 Thread 类中的 start()方法启动线程,执行 run 中的内容。 3.4.3 方法举例 public class TestRun implements Runnable { @Override public void run() { System.out.println("TestRun is running!"); }} public static void main(String[] args) { Thread thread = new Thread(new TestRun()); thread.start(); } 执行结果: 如结果所示:当线程实行 start 方法时,执行 Runnable 的 run 方法中的内容。 4 Function 接口用法 Function 的主要用途是可以通过 lambda 表达式实现方法的内容。 4.1 差异处理 原代码: @Datapublic class User { /** * 姓名 */ private String name; /** * 年龄 */ private int age; /** * 组员 */ private List<User> parters;} public static void main(String[] args) { User user =new User(); if(user ==null ||user.getAge() <18 ){ throw new RuntimeException("未成年!"); }} 执行结果: 使用 Function 接口后的代码: @FunctionalInterfacepublic interface testFunctionInfe { /** * 输入异常信息 * @param message */ void showExceptionMessage(String message);} public static testFunctionInfe doException(boolean flag){ return (message -> { if (flag){ throw new RuntimeException(message); } }); } public static void main(String[] args) { User user =new User(); doException(user ==null ||user.getAge() <18).showExceptionMessage("未成年!");} 执行结果: 使用 function 接口前后都抛出了指定的异常信息。 4.2 处理 if…else… 原代码: public static void main(String[] args) { User user =new User(); if(user==null){ System.out.println("新增用户"); }else { System.out.println("更新用户"); }} 使用 Function 接口后的代码: public static void main(String[] args) { User user =new User(); Consumer trueConsumer = o -> { System.out.println("新增用户"); }; Consumer falseConsumer= o -> { System.out.println("更新用户"); }; trueOrFalseMethdo(user).showExceptionMessage(trueConsumer,falseConsumer);}public static testFunctionInfe trueOrFalseMethdo(User user){ return ((trueConsumer, falseConsumer) -> { if(user==null){ trueConsumer.accept(user); }else { falseConsumer.accept(user); } });}@FunctionalInterfacepublic interface testFunctionInfe { /** * 不同分处理不同的事情 * @param trueConsumer * @param falseConsumer */ void showExceptionMessage(Consumer trueConsumer,Consumer falseConsumer);} 执行结果: 4.3 处理多个 if 原代码: public static void main(String[] args) { String flag=""; if("A".equals(flag)){ System.out.println("我是A"); }else if ("B".equals(flag)) { System.out.println("我是B"); }else if ("C".equals(flag)) { System.out.println("我是C"); }else { System.out.println("没有对应的指令"); }} 使用 Function 接口后的代码: public static void main(String[] args) { String flag="B"; Map<String, Runnable> map =initFunctionMap(); trueOrFalseMethdo(map.get(flag)==null).showExceptionMessage(()->{ System.out.println("没有相应指令"); },map.get(flag));}public static Map<String, Runnable> initFunctionMap(){ Map<String,Runnable> result = Maps.newHashMap(); result.put("A",()->{System.out.println("我是A");}); result.put("B",()->{System.out.println("我是B");}); result.put("C",()->{System.out.println("我是C");}); return result;}public static testFunctionInfe trueOrFalseMethdo(boolean flag){ return ((runnable, falseConsumer) -> { if(flag){ runnable.run(); }else { falseConsumer.run(); } });} 执行结果: 5 总结 Function 函数式接口是 java 8 新加入的特性,可以和 lambda 表达式完美结合,是非常重要的特性,可以极大的简化代码。

优秀的个人博客,低调大师

Autograd解析|OneFlow学习笔记

撰文|月踏 更新|赵露阳 前文《AI杂谈:手推BP》讲了Backward Propagation的数学原理。本文以OneFlow的代码为例,梳理Autograd模块的实现细节。 1 一个求梯度的小例子 先看下面这个简单的例子: import oneflow as ofx = of.randn(2, 2, requires_grad=True)y = x + 100z = y.sum()z.backward() forward pass可以对应到下面的计算图: 图1 即对应下面公式: 根据前文《 AI杂谈:手推‍BP 》很容易手动计算出x的梯度值,即: x1、x2、x3的计算过程类似,不再赘述,下面看一下OneFlow的执行结果,执行print(x.grad)可得到如下输出: tensor([[1., 1.], [1., 1.]], dtype=oneflow.float32) 可以看出,结果和前面公式(3)的计算结果一致,下面通过具体的代码实现来分析OneFlow的Autograd模块。 2 backward接口 上面例子中的python端的backward接口,调用的是python/oneflow/framework/tensor.py中的_backward接口: def _backward(self, gradient=None, retain_graph=False, create_graph=False): if not lazy_mode.is_enabled(): flow.autograd.backward(self, gradient, retain_graph, create_graph) else: ... 可以看到backward只支持eager模式,这是因为graph静态图模式下,计算图是提前编译好的,无需手动通过.backward()调用。flow.autograd.backward()会调用 oneflow/api/python/autograd/autograd.cpp 中导出的backward方法: ONEFLOW_API_PYBIND11_MODULE("autograd", m) { m.def("backward", &Backward); m.def("grad", &Grad);} 从pybind定义来看,这里面总共导出了两个接口(autograd.backward和autograd.grad)。其中,backward是对所有的requires_grad属性为True的节点求梯度,grad只对指定的叶子结点求梯度,原理上是相同的,本文只以backward为例来看代码的实现,backward接口会调用到同一个文件中的Backward函数: Maybe<one::TensorTuple> Backward(const one::TensorTuple& outputs, const one::TensorTuple& out_grads, bool retain_graph, bool create_graph) { if (create_graph) { retain_graph = true; } std::shared_ptr<one::TensorTuple> gradients = JUST(CheckAndInitOutGrads(outputs, out_grads)); JUST(one::GetThreadLocalAutogradEngine()->RunBackwardAndSaveGrads4LeafTensorIf( outputs, *gradients, retain_graph, create_graph)); return std::make_shared<one::TensorTuple>(0);} 这里的GetThreadLocalAutogradEngine()可以看作是一个thread_local的单例,位于 oneflow/core/autograd/autograd_engine.cpp ,返回一个autograd引擎(AutogradEngine)对象的指针: AutogradEngine* GetThreadLocalAutogradEngine() { thread_local static GraphAutogradEngine autograd_engine; return &autograd_engine;} AutogradEngine是OneFlow的Autograd的核心数据结构,它的继承关系如下: 图2 这里autograd引擎的子类实现有基于栈式的、基于图式的实现,默认使用基于图式的GraphAutogradEngine。从前面代码中可以看到,获取autograd引擎指针后,通过调用RunBackwardAndSaveGrads4LeafTensor函数,位于 oneflow/core/autograd/autograd_engine.cpp:L315 : Maybe<void> GraphAutogradEngine::RunBackwardAndSaveGrads4LeafTensor(const TensorTuple& outputs, const TensorTuple& out_grads, bool retain_graph, bool create_graph) { for (int i = 0; i < outputs.size(); ++i) { JUST(JUST(outputs.at(i)->current_grad())->PushPartialTensor(out_grads.at(i))); } GraphTask graph_task(outputs, retain_graph, create_graph); JUST(graph_task.ComputeDependencies()); JUST(graph_task.Apply(/*save_grad_for_leaf=*/true)); return Maybe<void>::Ok();} 这就真正进入了autograd模块的内部处理流程,后面继续分析。 3 FunctionNode和建立反向图 在进行backward pass时,执行的是一张反向图,反向图中的节点是在forward pass的时候建立的,其中的每个节点被称作FunctionNode,主要数据结构如下: 图3 先说图3中FunctionNode( oneflow/core/autograd/autograd_engine.h:L42) ,包含next_functions_、input_meta_data_、output_meta_data_这三个数据成员,其中next_functions_表示出边,另外两个表示一些meta信息,下面列几个主要的: is_leaf_:是不是叶子节点 requires_grad_:是不是需要求梯度值 retain_grad_:对于非叶子节点,是不是保存梯度值 acc_grad_:在gradient accumulation的的情况下,多个mini-batch的梯度累加 current_grad_:当前这个batch的梯度值 我们用到的是GraphFunctionNod( oneflow/core/autograd/autograd_engine.cpp:L178) GraphFunctionNode::GraphFunctionNode(const std::string& name, const std::shared_ptr<BackwardFunction>& backward_fn, const TensorTuple& inputs, const TensorTuple& outputs) : FunctionNode(name, backward_fn) { input_meta_data_.resize(inputs.size()); next_functions_.reserve(inputs.size()); for (int i = 0; i < inputs.size(); ++i) { if (inputs.at(i)->requires_grad()) { input_meta_data_.at(i) = inputs.at(i)->mut_autograd_meta(); next_functions_.emplace_back(inputs.at(i)->mut_grad_fn_node()); } } output_meta_data_.resize(outputs.size()); output_tensor_infos_.reserve(outputs.size()); for (int i = 0; i < outputs.size(); ++i) { const auto& autograd_meta = NewAutogradMeta(outputs.at(i)->requires_grad(), outputs.at(i)->is_leaf()); outputs.at(i)->set_autograd_meta(autograd_meta); output_meta_data_.at(i) = outputs.at(i)->mut_autograd_meta(); output_tensor_infos_.emplace_back(TensorInfo(*outputs.at(i))); } backward_fn_ = backward_fn;} 可见它主要对FunctionNode中的重要数据成员做了初始化,其中input_meta_data_、output_meta_data_中的AutogradMeta信息是从相应的input、output tensor中获取的,tensor通过桥接模式保存了一个TensorImpl对象指针,这个TensorImpl对象则维护了一个AutogradMeta对象。 继续看下FunctionNode中的反向函数backward_fn_,在《 OneFlow学习笔记:从Functor到OpExprInterpreter 》中讲到了在进行一个op调用的时候会执行AutogradInterpreter::Apply这个函数( oneflow/core/framework/op_interpreter/op_interpreter.cpp:L86 ),里面会创建这个反向函数: Maybe<void> AutogradInterpreter::Apply( const OpExpr& op_expr, const TensorTuple& inputs, TensorTuple* outputs, const OpExprInterpContext& ctx) const { ... autograd::AutoGradMode mode(false); JUST(internal_->Apply(op_expr, inputs, outputs, ctx)); std::shared_ptr<OpExprGradClosure> grad_closure(nullptr); if (requires_grad && !LazyMode::is_enabled()) { grad_closure = JUST(op_expr.GetOrCreateOpGradClosure()); auto backward_fn = std::make_shared<BackwardFunction>(); backward_fn->body = [=](const TensorTuple& out_grads, TensorTuple* in_grads, bool create_graph) -> Maybe<void> { autograd::AutoGradMode mode(create_graph); JUST(grad_closure->Apply(out_grads, in_grads)); return Maybe<void>::Ok(); }; backward_fn->status = [=]() { return grad_closure->state()->SavedTensors().size() > 0; }; JUST(GetThreadLocalAutogradEngine()->AddNode(op_expr.op_type_name() + "_backward", backward_fn, inputs, outputs)); } ... return Maybe<void>::Ok();} 可以看到反向图节点的名字是以正向图op的type name加上_backward的后缀来组成的,使用AddNode方法来创建FunctionNode( oneflow/core/autograd/autograd_engine.cpp:L356 ) Maybe<FunctionNode> GraphAutogradEngine::AddNode( const std::string& name, const std::shared_ptr<BackwardFunction>& backward_fn, const TensorTuple& inputs, TensorTuple* outputs) { // Firstly push function_node of tensor in stack which is leaf and requires_grad for (const std::shared_ptr<Tensor>& in_tensor : inputs) { if (in_tensor->is_leaf() && in_tensor->requires_grad()) { if (!in_tensor->grad_fn_node()) { JUST(AddAccumulateFunctionNode(in_tensor)); } } } std::shared_ptr<FunctionNode> func_node = std::make_shared<GraphFunctionNode>(name, backward_fn, inputs, *outputs); for (const std::shared_ptr<Tensor>& out_tensor : *outputs) { out_tensor->set_grad_fn_node(func_node); } return func_node;} 可见FunctionNode是挂在Tensor上的,通过Tensor的set_grad_fn_node接口维护到Tensor的数据结构中,在《 OneFlow学习笔记: Gl o bal View的相关概念和实现 》中画过Tensor的继承关系图,FunctionNode就是保存在TensorIf中: 图4 至此,已经理清了FunctionNode中各个成员的作用以及来历,假如以第二节的图1为例来画出对应的反向图的话,如下图所示: 图5 计算好的梯度值会被放到output_meta_data_中得AutogradMeta中,它可以通过tensor的acc_grad、current_grad接口来获取。 4 反向图的执行流程 接第三节列出的最后一段代码,其中最重要的两句话是: ...JUST(graph_task.ComputeDependencies());JUST(graph_task.Apply(/*save_grad_for_leaf=*/true));... 这里面的graph_task是GraphTask类型,它是一个很重要的数据结构,用来调度反向图中所有FunctionNode的执行,下面列一下它的主要成员: class GraphTask final { bool retain_graph_; bool create_graph_; std::vector<FunctionNode*> roots_; HashMap<FunctionNode*, int> dependencies_; HashSet<FunctionNode*> need_execute_;}; 先看本节开头的graph_task.ComputeDependencies,它主要是在初始化dependencies_这个map,这个map维护了每个FunctionNode的入度信息,再看graph_task.Apply,它主要是在通过拓扑序来访问反向图中的每个FunctionNode,并且对当前的FunctionNode进行各种操作( oneflow/core/autograd/autograd_engine.cpp:L287 ) Maybe<void>GraphTask::Apply(boolsave_grad_for_leaf){ std::queue<FunctionNode*> queue; for (FunctionNode* node : roots_) { if (dependencies_[node] == 0) { queue.push(node); } } while (!queue.empty()) { FunctionNode* node = queue.front(); queue.pop(); if (!need_execute_.empty() && need_execute_.find(node) == need_execute_.end()) { node->ReleaseOutTensorArgs(); continue; } if (/*bool not_ready_to_apply=*/!(JUST(node->Apply(create_graph_)))) { continue; } if (save_grad_for_leaf) { JUST(node->AccGrad4LeafTensor(create_graph_)); } JUST(node->AccGrad4RetainGradTensor()); node->ReleaseOutTensorArgs(); if (!retain_graph_) { node->ReleaseData(); } for (const auto& next_grad_fn : node->next_functions()) { FunctionNode* next_node = next_grad_fn.get(); dependencies_[next_node] -= 1; if (dependencies_[next_node] == 0) { queue.push(next_node); } } } return Maybe<void>::Ok();} 这里最重要的是下面两个语句: node->Apply node->AccGrad4LeafTensor 下面来逐个分析,先看node->Apply( oneflow/core/autograd/autograd_engine.cpp:L143 ),首先利用output_meta_data_初始化了output_grads,把它作为反向函数的输入,调用反向函数来求梯度值,求出的梯度值暂存在input_grads中,然后再更新到input_meta_data_中: Maybe<bool> FunctionNode::Apply(bool create_graph) { ... JUST(backward_fn_->body(output_grads, &input_grads, create_graph)); for (int i = 0; i < input_meta_data_.size(); ++i) { if (input_grads.at(i)) { ... JUST(input_meta_data_.at(i)->current_grad()->PushPartialTensor(input_grads.at(i))); } } return true;} 再看node->AccGrad4LeafTensor,这个函数最终会调用到CopyOrAccGrad,它主要用于在gradient accumulation的时候,多个mini-batch之间把梯度值多累加,和如果有hook函数的的话,使用注册的hook对当前的梯度值进行处理: Maybe<void> CopyOrAccGrad(AutogradMeta* autograd_meta, bool autograd_mode) { autograd::AutoGradMode mode(autograd_mode); auto current_grad = JUST(autograd_meta->current_grad()->GetAccTensor({})); if (!current_grad) { return Maybe<void>::Ok(); } if (autograd_meta->acc_grad()) { ... DevVmDepObjectConsumeModeGuard guard(DevVmDepObjectConsumeMode::NONE); const auto& output = JUST(functional::Add(autograd_meta->acc_grad(), current_grad, /*alpha=*/1, /*inplace=*/autograd_meta->is_grad_acc_inplace())); JUST(autograd_meta->set_acc_grad(output)); } else { JUST(autograd_meta->set_acc_grad(current_grad)); } for (const auto& hook : autograd_meta->post_grad_accumulation_hooks()) { auto new_grad = hook(autograd_meta->acc_grad()); if (new_grad) { JUST(autograd_meta->set_acc_grad(new_grad)); } } return Maybe<void>::Ok();} (特别感谢同事yinggang中间的各种答疑解惑。本文主要参考代码:https://github.com/Oneflow-Inc/oneflow/commit/a4144f9ecb7e85ad073a810c3359bce7bfeb05e1) 其他人都在看 25倍性能加速,OneFlow“超速”了 一个GitHub史上增长最快的AI项目 手把手推导Ring All-reduce的数学性质 DeepMind爆发史:决定AI高峰的“游戏玩家” 解读Pathways(二):向前一步是OneFlow 五年ML Infra生涯,我学到最重要的3个教训 OneFlow v0.7.0发布:全新分布式接口,LiBai、Serving等一应俱全 欢迎下载体验OneFlow v0.7.0:https://github.com/Oneflow-Inc/oneflow/ 本文分享自微信公众号 - OneFlow(OneFlowTechnology)。 如有侵权,请联系 support@oschina.cn 删除。 本文参与“OSC源创计划”,欢迎正在阅读的你也加入,一起分享。

优秀的个人博客,低调大师

Greenplum测试框架最全解析

软件测试是开发过程中十分重要的一环,在数据库领域更是如此。一款稳定、可靠的数据库离不开大量的测试作为支撑。 Greenplum 作为一款基于 Postgres 的开源数据库,在测试方面做出了大量的探索。除继承了 Postgres 原有的 regress 测试外,增加了 Fault Injector 框架。允许开发者在回归测试中,通过执行简单的 SQL 函数,对数据库注入真实场景中可能出现各种的故障。此外, Greenplum 还开发了新的 isolation 测试框架 (isolation2),开发者可以用更简单易懂的语法编写出在各种并发情况下,可能出现的数据竞争的测试用例。配合 Fault Injector 框架,能够涵盖范围非常广的测试场景。 本文主要结合一些实际的案例介绍这些测试框架的使用场景和使用方式, 在后续的文章中会介绍这些框架的原理,欢迎大家留言交流。 regress regress 测试位于 Greenplum/Postgres 源代码的 src/test/regress 目录下。这些测试通常包含: 测试用例:一组 SQL 语句 (.sql 文件) 期待输出:测试用例执行后, Postgres 的正确输出 (.out 文件) 这些测试由 pg_regress 调度执行。执行结束后通过对预期输出和实际输出进行对比 (通常会自动生成一个 regression.diffs 的文件),即可知道哪些测试没有通过,以及原因是什么。利用 regress 测试,我们可以做一些功能性的测试。例如:测试一些带有复杂过滤条件的 SELECT 语句的输出结果是否正确,测试一些 SQL 生成的执行计划是否符合预期。这些测试都有一个特点,那就是不管在什么情况下,它们的输出都应该是一致且恒定的。 尽管 regress 测试能够满足大部分的功能性测试,但还有一些有趣的测试场景值得讨论。比如数据库中的隔离和并发问题。虽然 pg_regress 支持并行地执行多个测试用例,但我们并不能利用它来验证一些时序上的问题 (我们倒是可以利用它来加速一些测试)。因为这些并行的测试用例中 SQL 语句的执行次序是没有保证的,我们无法得到稳定的输出,也就无法跟预期的输出进行对比了。而下面要介绍的 isolation & isolation2 就是为这类测试而设计的。 isolation & isolation2 isolation 测试位于 Greenplum/Postgres 源代码的 src/test/isolation 目录下。目前 Greenplum 中的 isolation 测试框架还不是很完善,Greenplum 仅在 Utility 模式下运行 Postgres 原生的 isolation 测试。Greenplum 大部分的和并发/隔离相关的测试由 isolation2 完成。 Postgres 原生的 isolation 测试包含: 测试用例:一组 .spec 文件,这些 .spec 文件除了包含待测试的 SQL 语句外,最重要的是还包含了对执行这些 SQL 语句顺序的定义。 期待输出:测试用例的 spec 被执行后, Postgres 的正确输出 (.out 文件)。 这些测试由 pg_isolation_regress 根据 .spec 文件中定义的执行顺序,调度不同的 Session 执行,下面给出的是 Postgres 中, 对于一种死锁的测试用例。 ## file: src/test/isolation/specs/deadlock-simple.spec setup { CREATE TABLE a1 (); } teardown { DROP TABLE a1; } session s1 setup { BEGIN; } step s1as { LOCK TABLE a1 IN ACCESS SHARE MODE; } step s1ae { LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; } step s1c { COMMIT; } session s2 setup { BEGIN; } step s2as { LOCK TABLE a1 IN ACCESS SHARE MODE; } step s2ae { LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; } step s2c { COMMIT; } permutations1ass2ass1aes2aes1cs2c 其中,setup 和 teardown 中包含的 SQL 语句分别会在测试运行前和结束后运行,在这个例子中是创建和删除表 a1. s1 和 s2 分别是这个测试中两个 Session 的名称。 在 s1 中,定义了 setup 时需要开始一段事务,还定义了 s1as, s1ae, s1c 这三个执行步骤: s1as:对表 a1 以 ACCESS SHARE 模式上锁 s1ae:对表 a1 以 ACCESS EXCLUSIVE 模式上锁 s1c:提交当前事务 在 s2 中,定义了 setup 时需要开始一段事务,还定义了 s2as, s2ae, s2c 这三个执行步骤: s2as: 对表 a1 以 ACCESS SHARE 模式上锁 s2ae: 对表 a1 以 ACCESS EXCLUSIVE 模式上锁 s2c: 提交当前事务 permutation 中定义了这些 SQL 语句执行的顺序: s1: 对表 a1 以 ACCESS SHARE 模式上锁 s2: 对表 a1 以 ACCESS SHARE 模式上锁 s1: 对表 a1 以 ACCESS EXCLUSIVE 模式上锁 <等待, s2 拿了 ACCESS SHARE 模式的锁> s2: 对表 a1 以 ACCESS EXCLUSIVE 模式上锁 <死锁发生, s1, s2 互相等待> pg_isolation_regress 执行完上面 .spec 文件中定义的 SQL 后, 正确输出如下: ## file: src/test/isolation/expected/deadlock-simple.out Parsed test spec with 2 sessions starting permutation: s1as s2as s1ae s2ae s1c s2c step s1as: LOCK TABLE a1 IN ACCESS SHARE MODE; step s2as: LOCK TABLE a1 IN ACCESS SHARE MODE; step s1ae: LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; <waiting ...> step s2ae: LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; ERROR: deadlock detected step s1ae: <... completed> step s1c: COMMIT; steps2c:COMMIT; 开发者可以在 .spec 文件中灵活的描述在不同 Session 中执行 SQL 语句的顺序,使得一些与并发/隔离相关的测试更为容易编写。 Greenplum 团队开发的 isolation2 测试框架,也是类似的思路,但 isolation2 的语法更简单易懂。下面是笔者利用 isolation2 编写的与上面一致的测试用例。 CREATE TABLE a1(); 1: BEGIN; -- s1 开始一段事务 1: LOCK TABLE a1 IN ACCESS SHARE MODE; -- s1 以 ACCESS SHARE 模式对 a1 上锁 2: BEGIN; -- s2 开始一段事务 2: LOCK TABLE a1 IN ACCESS SHARE MODE; -- s2 以 ACCESS SHARE 模式对 a1 上锁 1&: LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; -- s1 以 ACCESS EXCLUSIVE 模式对 a1 上锁 2: LOCK TABLE a1 IN ACCESS EXCLUSIVE MODE; -- s2 以 ACCESS EXCLUSIVE 模式对 a1 上锁 1<: -- 等待 s1 返回 1: COMMIT; -- s1 提交事务 2: COMMIT; -- s2 提交事务 DROPTABLEa1; 其中不同的数字表示在不同的 Session 中执行对应的 SQL,由上至下表示了 SQL 的执行顺序。‘&’ 表示当前的 Session 执行的 SQL 会被阻塞,‘<’ 表示等待当前的 Session 返回。isolation2 支持的语法非常丰富,这里就不一一列举了,有兴趣的读者可以浏览 Greenplum 源代码中 src/test/isolation2/sql 目录下的测试用例。 通过对比 isolation 与 isolation2,在功能方面,二者并没有很大的区别;在语法方面, isolation 的 .spec 文件更为严谨,当测试用例中需要多次调用相同的 SQL 语句时会十分方便。solation2 的语法更为灵活直观, 编写起来比较容易上手。 FaultInjector 真实世界中的故障往往要复杂许多,比如:网络可能会有很大的延迟,磁盘中的数据可能丢失,集群中的某台主机可能突然下线。Greenplum 团队开发的 Fault Injector 框架让这类测试变得十分容易。例如,我们希望在 heap 中插入 tuple 的时候,注入一些故障,观察这些故障对系统的影响。在 Greenplum master 分支的 src/backend/access/heap/heapam.c 中,有如下代码: // 注: 笔者为了叙述方便, 对这里的代码进行了调整, 与实际代码可能有些出入. void heap_insert(Relation relation, HeapTuple tup, CommandId cid, int options, BulkInsertState bistate, TransactionId xid) { ... #ifdef FAULT_INJECTOR FaultInjector_InjectFaultIfSet("heap_insert", /*faultname*/ DDLNotSpecified, /*ddlstatement*/ "", /*databasename*/ RelationGetRelationName(relation)); #endif ... } 当我们启用了 Fault Injector 后,每当代码执行到 FaultInjector_InjectFaultIfSet() 时,Fault Injector 会检测当前的注入点是否被注入了故障,如果有,则执行相关的故障逻辑。利用 gp_inject_fault 插件可以很容易的在注入点注入故障逻辑。下面这段 SQL 演示了如何在 Greenplum 集群中,在编号为 1 的 Segment Server 上, ‘heap_insert’ 这一注入点,注入一段死循环。以此来模拟实际场景中,某个 Segment Server 发生了故障,无法正常地插入 tuple 到表 my_table 中的情况。 CREATE EXTENSION gp_inject_fault; SELECT gp_inject_fault('heap_insert' /* inject location */, 'infinite_loop' /* fault type */, '' /* DDL */, '' /* database name */, 'my_table' /* table name */, 1 /* start occurrence */, 10 /* end occurrence */, 0 /* extra arg */, dbid) FROMgp_segment_configurationWHEREcontent=1ANDrole='p'; 除了可以注入 ‘infinite_loop’ 外,Fault Injector 还支持注入以下常见的几种故障: error:等价于 elog(ERROR) fatal:等价于 elog(FATAL) panic:等价于 elog(PANIC) sleep:休眠一段时间 suspend:阻塞当前的进程, 并且不检查中断信号 resume:恢复被注入 ‘suspend’ 故障的进程 skip:用来注入一些自定义的故障逻辑 reset:移除之前注入的故障 segv:使当前的进程崩溃 (发送 SIGSEGV 信号) ... 有关其它的故障类型以及更多的使用方法,可以参阅 gpcontrib/gp_inject_fault/README。目前 Postgres 还不支持 Fault Injector 框架, Greenplum 团队向上游提交了 Patch,感兴趣的读者可以下载这个 Patch[1] 体验。 后 记 ‍ ‍ ‍ ‍ ‍ ‍ 笔者自己上手一个新项目时,喜欢从一些测试用例开始摸索。通过一些简单的测试用例可以窥知一款软件的基本功能,之后配合 git blame 便可以顺藤摸瓜,找到引入这些测试用例的 commit 记录。有时候它可能是修复了一个 Bug,也可能是引入了一个新功能,这样就可以学习到开发这个功能时该如何写测试,该去修改哪些代码。最后,希望大家多多为 Greenplum 贡献代码/Issue。 参考资料 [1] Fault Injector Framework: https://www.postgresql.org/message-id/flat/CANXE4TdxdESX1jKw48xet-5GvBFVSq%3D4cgNeioTQff372KO45A%40mail.gmail.com 郭兴,VMware Greenplum实习生 东南大学研究生在读, 2021 年 9 月加入 Greenplum 团队开启实习工作, 参与 Greenplum Extension 研发。 点击文末“ 阅读原文 ”,获取Greenplum中文资源。 来一波 “在看”、“分享” 和 “赞” 吧! 本文分享自微信公众号 - Greenplum中文社区(GreenplumCommunity)。如有侵权,请联系 support@oschina.cn 删除。本文参与“OSC源创计划”,欢迎正在阅读的你也加入,一起分享。

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册