首页 文章 精选 留言 我的

精选列表

搜索[源码学习],共10000篇文章
优秀的个人博客,低调大师

ElasticSearch Bulk 源码解析

本来应该先有这篇文章,后有 如何提高ElasticSearch 索引速度才对。不过当时觉得后面一篇文章会更有实际意义一些,所以先写了后面那篇文章。结果现在这篇文章晚了20多天。 前言 读这篇文章前,建议先看看ElasticSearch Rest/RPC 接口解析,有利于你把握ElasticSearch接受处理请求的脉络。对于RPC类的调用,我会在后文简单提及,只是endpoint不一样,内部处理逻辑还是一样的。这篇只会讲IndexRequest,其他如DeleteRequest,UpdateRequest之类的,我们暂时不涉及。 类处理路径 RestBulkAction -> TransportBulkAction -> TransportShardBulkAction 其中TransportShardBulkAction比较特殊,有个继承结构: TransportShardBulkAction < TransportReplicationAction < TransportAction 主入口是TransportAction,具体的业务逻辑实现分布到子类(TransportReplicationAction)和孙子类(TransportShardBulkAction)里了。 另外,我们也会提及org.elasticsearch.index.engine.Engine相关的东西,从而让大家清楚的了解ES是如何和Lucene关联上的。 RestBulkAction 入口自然是org.elasticsearch.rest.action.bulk.RestBulkAction,一个请求会构建一个BulkRequest对象,BulkRequest.add方法会解析你提交的文本。对于类型为index或者create的(还记得bulk提交的文本格式是啥样子的么?),都会被构建出IndexRequest对象,这些解析后的对象会被放到BulkRequest对象的属性requests里。当然如果是update,delete等则会构建出其他对象,但都会放到requests里。 public class BulkRequest extends ActionRequest<BulkRequest> implements CompositeIndicesRequest { //这个就是前面提到的requests final List<ActionRequest> requests = new ArrayList<>(); //这个复杂的方法就是通过http请求参数解析出 //IndexRequest,DeleteRequest,UpdateRequest等然后放到requests里 public BulkRequest add(BytesReference data, @Nullable String defaultIndex, @Nullable String defaultType, @Nullable String defaultRouting, @Nullable String[] defaultFields, @Nullable Object payload, boolean allowExplicitIndex) throws Exception { XContent xContent = XContentFactory.xContent(data); int line = 0; int from = 0; int length = data.length(); byte marker = xContent.streamSeparator(); while (true) { 接着通过NodeClient将请求发送到TransportBulkAction类(回忆下之前文章里提到的映射关系,譬如 TransportAction,两层映射关系解析 )。对应的方法如下: //这里的client其实是NodeClient client.bulk(bulkRequest, new RestBuilderListener<BulkResponse>(channel) { TransportBulkAction 看这个类的签名: public class TransportBulkAction extends HandledTransportAction<BulkRequest, BulkResponse> { 实现了HandledTransportAction,说明这个类同时也是RPC接口的逻辑处理类。如果你点进HandledTransportAction就能看到ES里经典的messageReceived方法了。这个是题外话 该类对应的入口是: protected void doExecute(final BulkRequest bulkRequest, final ActionListener<BulkResponse> listener) {这里的bulkRequest 就是前面RestBulkAction组装好的。该方法第一步是判断是不是需要自动建索引,如果索引不存在,就自动创建了。 接着通过executeBulk方法进入原来的流程。在该方法中,对bulkRequest.requests 进行了两次for循环。 第一次判定如果是IndexRequest就调用IndexRequest.process方法,主要是为了解析出timestamp,routing,id,parent 等字段。 第二次是为了对数据进行分拣。大致是为了形成这么一种结构: //这里的BulkItemRequest来源于 IndexRequest等 Map[ShardId, List[BulkItemRequest]] 接着对新形成的这个结构(ShardId -> List[BulkItemRequest])做循环,也就是针对每个ShardId里的数据进行统一处理。有了ShardId,bulkRequest,List[BulkItemRequest]等信息后,统一封装成BulkShardRequest。从名字看就很好理解,就是对属于同一ShardId的数据构建一个新的类似BulkRequest的对象。 接着就到TransportShardBulkAction,TransportReplicationAction,TransportAction 三代人出场了: //这里的shardBulkAction 是TransportShardBulkAction shardBulkAction.execute(bulkShardRequest, new ActionListener<BulkShardResponse>() { TransportReplicationAction/TransportShardBulkAction TransportAction是一个通用的主类,具体逻辑还是其子类来实现。虽然前面提到shardBulkAction是TransportShardBulkAction,但其实流程逻辑还是TransportReplicationAction来完成的。入口在该类的doExecute方法: @Override protected void doExecute(Request request, ActionListener<Response> listener) { new PrimaryPhase(request, listener).run(); } 我们知道在ES里有主从分片的概念,所以一条数据被索引后需要经过两个阶段: 将数据写入Primary(主分片) 将数据写入Replication(从分片) 至于为什么不直接从Primary进行复制,而是将数据分别写入到Primary和Replication我觉得主要考虑如果一旦Primary是损坏的,不至于影响到Replication(考虑下,如果Primary是损坏的文件,然后所有的Replication如果是直接复制过来,就都坏了)。 又扯远了。我们看到doExecute 首先是进入PrimaryPhase阶段,也就是写主分片。 Primary Phase 在PrimaryPhase.doRun方法里,你会看到两行代码 final ShardIterator shardIt = shards(observer.observedState(), internalRequest); final ShardRouting primary = resolvePrimary(shardIt); 其中这个ShardIterator是类似 shardId->ShardGroup 的结构。不管这个shardId是什么,它一定是个Replication或者Primary的shardId, ShardGroup 就是Replication和Primary的集合。resolvePrimary方法则是遍历这个集合,然后找出Primary的过程。 知道Primary后就可以判断是转发到别的Node或者直接在本Node处理了: routeRequestOrPerformLocally(primary, shardIt); 如果Primary就在本节点,直接就处理了: //我去掉了一些无关代码哈 if (primary.currentNodeId().equals(observer.observedState().nodes().localNodeId())) { try { threadPool.executor(executor).execute(new AbstractRunnable() { @Override protected void doRun() throws Exception { performOnPrimary(primary, shardsIt); } } 这里用上了线程池。前面对每个shardId对应的数据集合做处理,其实是顺序循环执行的,这里实现了将数据处理异步化。 在performOnPrimary方法中,BulkShardRequest被转化成了PrimaryOperationRequest,理由也很简单,更加specific了,因为就是针对主分片的Request。接着进入shardOperationOnPrimary 方法,该方法是在孙子类TransportShardBulkAction类里实现的。 protected Tuple<BulkShardResponse, BulkShardRequest> shardOperationOnPrimary( ClusterState clusterState, PrimaryOperationRequest shardRequest) { 到该方法,有两个比较重要的概念会出现: //伟大的版本号,实现了对并发修改的支持 long[] preVersions = new long[request.items().length]; VersionType[] preVersionTypes = new VersionType[request.items().length]; //事物日志,为Shard Recovery以及 //避免过多的Index Commit做出突出贡献, //同时也是是实现了GetById的实时性 Translog.Location location = null; 上面两个概念成就了ES从一个简单的全文检索引擎到类No-SQL的转型(好吧,我好像又扯远了) 接着就是for循环了: //这里的request是BulkShardRequest //对应的items则是BulkItemRequest集合 for (int requestIndex = 0; requestIndex < request.items().length; requestIndex++) { 循环会根据BulkItemRequest的不同类型而有了分支。其实就是 IndexRequest,DeleteRequest,UpdateRequest,我们这里依然只讨论IndexRequest。如果发现BulkItemRequest是IndexRequest,进行如下操作: WriteResult<IndexResponse> result = shardIndexOperation(request, indexRequest, clusterState, indexShard, true); shardIndexOperation里嵌套的核心方法是executeIndexRequestOnPrimary,该方法第一步是获取到Operation对象, Engine.IndexingOperation operation = prepareIndexOperationOnPrimary(shardRequest, request, indexShard); Engine对象是比较底层的一个对象了,是对Lucene的IndexWriter,Searcher之类的封装。这里的Engine.IndexingOperation对应的是Create或者Index类。你可以把这两个类理解为待索引的Document,只是还带上了动作。 第二步是判断索引的Mapping是不是要动态更新,如果是,则更新。 第三步执行实际的建索引操作: final boolean created = operation.execute(indexShard); operation.execute 额外引出的话题 我们会暂时深入到operate.execute方法里,但这个不是主线,看完后记得回到上面那行代码上。 刚才我们说了operation可能是Create或者Index,我们会以Create为主线进行分析。所谓Create和Index,你可以理解为一个待索引的Document,只是带上动作的语义。 上面对应的execute 方法签名是: @Overridepublic boolean execute(IndexShard shard) { shard.create(this); return true; } 我们看到这里是反向调用indexShard对象的create方法来进行索引的创建。我们来看看IndexShard的create方法: //我依然做了删减,体现一些核心代码 public void create(Engine.Create create) { engine().create(create); } engine()方法返回的是InternalEngine实例,InternalEngine .innerCreate方法执行到构建索引的操作。这个方法值得分析一下,所以我就贴了一坨的代码。 private void innerCreate(Create create) throws IOException { if (engineConfig.isOptimizeAutoGenerateId() && create.autoGeneratedId() && !create.canHaveDuplicates()) { // We don't need to lock because this ID cannot be concurrently updated: innerCreateNoLock(create, Versions.NOT_FOUND, null); } else { synchronized (dirtyLock(create.uid())) { final long currentVersion; final VersionValue versionValue; versionValue = versionMap.getUnderLock(create.uid().bytes()); if (versionValue == null) { currentVersion = loadCurrentVersionFromIndex(create.uid()); } else { if (engineConfig.isEnableGcDeletes() && versionValue.delete() && (engineConfig.getThreadPool().estimatedTimeInMillis() - versionValue.time()) > engineConfig.getGcDeletesInMillis()) { currentVersion = Versions.NOT_FOUND; // deleted, and GC } else { currentVersion = versionValue.version(); } } innerCreateNoLock(create, currentVersion, versionValue); } } } 首先,如果满足如下三个条件就无需进行版本检查: index.optimize_auto_generated_id 被设置为true(默认是false,话说注释上说是默认是true,但是我看着觉得像是false) id设置为自动生成(没有人工设置id) create.canHaveDuplicates == false ,该参数一般是false 提这个是主要为了说明,譬如一般的运维日志啥的,就不要自己生成ID了,采用自动生成的ID,可以跳过版本检查,从而提高入库的效率。 第二个指的说的是,如果对应文档在缓存中没有找到(versionMap),那么就会由如下的代码执行实际磁盘查询操作: currentVersion = loadCurrentVersionFromIndex(create.uid()); 通过对比create对象里的版本号和从索引文件里加载的版本号 ,最终决定是进行update还是create操作。 在innerCreateNoLock 方法里,你会看到熟悉的Lucene操作,譬如: indexWriter.addDocument(index.docs().get(0)); //或者 indexWriter.updateDocument(index.uid(), index.docs().get(0)); 现在回到TransportShardBulkAction的主线上。执行完下面的代码后: final boolean created = operation.execute(indexShard); 就能获得对应文档的版本等信息,这些信息会更新对应的IndexRequest等对象。 到目前为止,Primay Phase 完成,接着开始Replication Phase replicationPhase = new ReplicationPhase(shardsIt, primaryResponse.v2(), primaryResponse.v1(), observer, primary, internalRequest, listener, indexShardReference); finishAndMoveToReplication(replicationPhase); 最后一行代码会启动replicationPhase阶段。 Replication Phase Replication Phase 流程大致和Primary Phase 相同,就不做过详细的解决,我这里简单提及一下。 ReplicationPhase的doRun方法是入口,核心方法是performOnReplica,如果发现Replication shardId所属的节点就是自己的话,异步执行shardOperationOnReplica,大体逻辑如下: threadPool.executor(executor).execute(new AbstractRunnable() { @Override protected void doRun() { try { shardOperationOnReplica(shard.shardId(), replicaRequest); onReplicaSuccess(); } catch (Throwable e) { onReplicaFailure(nodeId, e); failReplicaIfNeeded(shard.index(), shard.id(), e); } } 在Replication阶段,shardOperationOnReplica 该方法完成了索引内容解析,mapping动态新增,最后进入索引(和就是前面提到的operation.execute)等动作,所以还是比Primary 阶段更紧凑些。 另外,在Primary Phase 和 Replication Phase, 一个BulkShardRequest 处理完成后(也就是一个Shard 对应的数据集合)才会刷写Translog日志。所以如果发生数据丢失,则可能是多条数据。 总结 这篇文章以流程分析为主,很多细节我们依然没有讲解详细,比如Translog和Version。这些争取能够在后续文章中进一步阐述。另外错误之处在所难免,请大家在评论处提出。

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

HBase源码阅读资源

HBase MemStoreFlusher 虽与最新版0.98.7的实现已经有差异,但分析的比较好 MemeStoreFlusher在HRegionServer类中初始化。 HRegionServer实现了Runnable接口,在run方法中针对MemeStoreFlusher进行了初始化 privatevoidinitializeThreads()throwsIOException{ //Cacheflushingthread. this.cacheFlusher=newMemStoreFlusher(conf,this); ... } 启动: this.cacheFlusher.start(uncaughtExceptionHandler); interrupt: if(this.cacheFlusher!=null)this.cacheFlusher.interruptIfNecessary();

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

短视频APP源码直播APP源码什么样的好

短视频所面临的架构问题短视频相比于文本数据而言,有着一些差异:1.数据大小的差异。比如一条美拍,经过视频压缩和清晰度的权衡,10s的视频大小1MB多,而一条5分钟视频的美拍甚至要达到几十M,相比与几十字节或者几百字节的文本要大得多。因为数据量要大得多,所以也会面临一些问题:如何上传、如何存放、以及如何播放的问题。关于上传,要在手机上传这么一个视频,特别是弱网环境要上传这么一个文件,上传的成功率会比较低,晚高峰的时候,省际网络的拥塞情况下,要更为明显得多。所以针对上传,需要基于CDN走动态加速来优化网络链路(通过基调实测过对于提升稳定性和速度有一定帮助),同时对于比较大的视频需要做好分片上传,减少失败重传的成本和失败概率等来提升可用性。同时不同CDN厂商的链路状况在不同的运营商不同地区可能表现不一,所以也需要结合基调测试,选择一些比较适合自己的CDN厂商链路。同时因为数据相对比较大,当数据量达到一定规模,存储容量会面临一些挑战,目前美拍的视频容量级别也达到PB级别的规模,所以要求存储本身能够具备比较强的线性扩展能力,并且有足够的资源冗余,而传统的Mysql等数据库比较难来支持这个场景,所以往往借助于专用的分布式对象存储,通过自建的服务或者云存储服务能够解决,得益于近几年云存储的发展,目前美拍主要还是使用云存储服务来解决。自身的分布式对象存储主要用于解决一些内部场景,比如对于数据隐私性和安全性要求比较高的场景。关于对于播放,因为文件比较大,也容易受到网络的影响,所以为了规避卡顿,一些细节也需要处理。比如对于60s,300s的视频,需要考虑到文件比较大,同时有拖动的需求,所以一般使用http range的方式,或者基于HLS的点播播放方式,基于前者比较简单粗暴,不过基于播放器的机制,也能够满足需求,也能实现点播拖动。而直接基于HLS的方式会更友好,特别是更长的一些视频,比如5分钟甚至更大的视频,不过这种需要单独的转码支持。之前美拍主要是短视频为主,所以更多使用http range的方式。而后续随着5分钟或者更大视频的场景,也在逐步做一些尝试。同时对于播放而言,在弱化环境下,可能也会面临一些问题,比如播放时长卡顿的问题,这种一般通过网络链路优化;或者通过多码率的 自适应优化,比如多路转码,然后根据特定算法模型量化用户网络情况进行选码率,网络差的用低码率的方式。2.数据的格式标准差异相比与文本数据,短视频本身是二进制数据,有比较固定的编码标准,比如H.264、H.265等,有着比较固定和通用的一些格式标准。3.数据的处理需求视频本身能够承载的信息比较多,所以会面临有大量的数据处理需求,比如水印、帧缩略图、转码等,以及短视频鉴黄等。而视频处理的操作是非常慢的,会带来巨大的资源开销。美拍对于视频的处理,主要分为两块:客户端处理,视频处理尽量往客户端靠,利用现有强大的手机处理性能来规避减少服务器压力,同时这也会面临一些低端机型的处理效率问题,不过特别低端的机型用于上传美拍本身比较少数,所以问题不算明显。客户端主要是对于视频的效果叠加、人脸识别和各种美颜美化算法的处理,我们这边客户端有实验室团队,在专门做这种效果算法的优化工作。同时客户端处理还会增加一些必要的转码和水印的视频处理。目前客户端的视频编解码方式,会有软编码和硬编码的方式,软编码主要是兼容性比较好,编码效果好些,不过缺点就是能耗高且慢些。而硬编码借助于显卡等,能够得到比较低的能耗并且更快,不过兼容和效果要差一些,特别是对于一些低配的机型。所以目前往往采用结合的方式。服务端的处理,主要是进行视频的一些审核转码工作,也有一些抽帧生成截图的工作等,目前使用ffmpeg进行一些处理。服务端本身需要考虑的一些点,就是因为资源消耗比较高,所以需要机器数会多,所以在服务端做的视频处理操作,会尽量控制在一个合理的范围。同时因为可能美拍这种场景,也会遇到这些热点事件的突变峰值,所以转码服务集群本身需要具备可弹性伸缩和异步化消峰机制,以便来适应这种突增请求的场景。 审核问题视频内容本身可以有任意多样的表现形式,所以也是一个涉黄涉恐的多发地带,而这是一个无法规避掉的需求,因为没有处理好,可能分分钟被封站。审核的最大的问题,主要是会面临视频时长过长,会带来人力审核成本的提升。比如100万个视频,每个平均是30s的话,那么就3000W 秒,大概需要347人日。 通过技术手段可以做一些工作,比如:接入一些比较好的第三方的视频识别模块,如果能够过滤掉85%保证没有问题的视频的话,那么工作量会缩减到15%。不过之前在接入使用的时候,发现效果没有达到预期,目前也在逐步尝试些其他方案。通过抽帧的方式,比如只抽取某几帧的方式进行检查。通过转码的方式,比如一个60s的美拍视频,通过2倍速的方式,无声,140 * 140的分辨率转换,大概大小能够在650kB左右,这样加速了播放的过程的同时,还能够减少审核带宽的消耗,减少了下载过程。基于大数据分析,分析一些高危地带、用户画像等,然后通过一些黑名单进行一些处理,或者对于某些潜在高危用户进行完整视频的审核,而对于低危用户进行抽帧的方式等等。

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

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

上一篇我们详细讲述“4.2.4 ustore”相关内容。本篇我们将继续为小伙伴们带来“4.2.5 行存储索引机制”、“4.2.6 行存储缓存机制”及“4.2.7 cstore”等精彩内容的详细介绍。 4.2.5 行存储索引机制 本节以B-Tree索引为例,介绍openGauss中行存储(格式)表的索引机制。索引本质上是对数据的一种物理有序聚簇。有序聚簇参考的排序字段被称为索引键。为了节省存储空间,一般索引表中只存储有序聚簇的索引键键值以及对应元组在主表中的物理位置。在查询指定的索引键键值元组时,得益于有序聚簇排序,可以快速找到目标元组在主表中的物理位置,然后通过访问主表对应页面和偏移得到目标元组。B-Tree索引的组织结构如图4-19所示。 图4-19 B-Tree索引页面间和页面内结构示意图 当前openGauss版本中,每个B-Tree的页面采用和行存储astore堆表页面基本相同的页面结构(见“4.2.3 astore”节的“2. astore堆表页面元组结构”小节)。页面间按照树形结构组织,分为根节点页面、内部节点页面和叶子节点页面。其中,根节点页面和内部节点页面中的索引元组不直接指向堆表元组,而是指向下一层的内部节点页面或叶子节点页面;叶子节点页面位于B-Tree的最底层,叶子节点页面中的索引元组指向索引键值对应的堆表元组,即存储了该元组在堆表中的物理位置(堆表页面号和页内偏移)。 B-Tree索引元组结构由索引元组头、NULL值字典和索引键值字段3分组成。 索引元组头为IndexTupleData结构体,定义代码如下所示。其中,t_tid为堆表元组的位置或下一层索引页面的位置;t_info为标志位,记录键值中是否有NULL值、是否有变长键值、索引访存方式信息以及元组长度。 typedef struct IndexTupleData { ItemPointerData t_tid; /* 堆表元组的物理行号 */ /* --------------- * t_info标志位内容: * * 第15位:是否有NULL字段 * 第14位:是否有变长字段 * 第13位: 访存方式自定义 * 第0-12位: 元组长度 * --------------- */ unsigned short t_info; /* 如上 */ } IndexTupleData; /* 实际索引元组数据紧跟该结构体 */ 与astore堆表元组不同,索引表的NULL值字典是定长的,一个bit位对应一个索引字段。当前最多支持32个索引字段,因此该字典的长度为4个字节(如果要支持变长,那么长度加变长字典的实际空间并不会比定长的4个字节少多少)。如果索引元组头部t_info标志位中存在NULL值的bit位为0,那么该索引元组没有NULL值字典,可以节约4个字节的空间。 索引键值字段和astore堆表元组的字段结构是完全相同的,唯一区别是索引键值只保存创建索引的那些字段上的值。 为了在一个索引页面中能够保存尽可能多的元组个数,降低整个B-Tree结构的层数,索引元组和astore堆表元组的结构相比要紧凑很多,去掉了一些和astore堆表元组冗余的结构体成员。在实际执行索引查询的时候,一般需要加载(索引层数+1)个物理页面才能找到目标元组。一般索引层数在2至4层之间,因此每减少一个层级近似就可以节省20%以上的元组访存开销。 当前openGauss版本中,索引元组头部不保存t_xmin和t_xmax这两个事务信息,因此元组可见性的判断不会在遍历索引时确定,而是要等到获得叶子索引最终指向的堆表元组以后,通过结合查询快照和堆表元组的t_xmin、t_xmax信息,才能判断对应堆表元组对本查询是否可见。将导致以下几个现象: 对于被删除的astore堆表元组,其空间(至少其元组指针)不能立刻被释放,否则会留下悬空的索引指针,导致后续查询出现问题。 对于被更新的astore堆表元组,如果更新前后索引字段的值发生变化,那么需要插入一条新的索引元组来指向更新后的堆表元组。然而即使更新前后所有索引字段的值没有发生变化,考虑到可能还有并发的查询需要访问老元组,因此老索引元组还要保留。同时要么插入一条新索引元组来指向更新后的堆表元组。或者也可以通过将更新后元组的位置信息保存在老元组中,这样通过原来的一条索引元组,就可以一并查到更新前后的两条新、老元组了。但是这种场景下老堆表元组的清理又变得复杂起来,否则还会存在悬挂索引指针的问题。 为了解决上述这些问题,openGauss当前提供了三种空间管理和回收的机制(参见“4.2.3 astore”节的“5. astore空间管理和回收”小节)。在对astore堆表进行轻量级清理时,无法清理索引中的垃圾数据。只有对astore进行中量级VACUUM清理,或者重量级VACUUM FULL清理时,才能够清理对应索引中的垃圾数据。 最后,上述索引可见性判断机制有一种例外场景:如果查询不涉及非索引字段,如显示查询索引字段内容、或“SELECT COUNT(*)”类查询,且索引字段t_tid指向的astore堆表页面对应的VM(visibility map,可见性位图)比特位为1,那么该索引元组被认为是可见的,这种扫描方式称为“Index Only Scan”。该扫描方式不仅提高了可见性判断的效率,更重要的是避免了对于堆表页面的访问,从而可以节省大量I/O开销。在页面空闲空间回收过程中,如果被清理的堆表页面上的所有元组对于当前所有正在执行的事务都可见,那么其对应的VM比特位会被置为1;后续如果该堆表页面上有新的插入、删除或更新操作之后,都会将其对应的VM比特位置为0。 openGauss中的行存储索引表访存接口如表4-26所示。 表4-26 行存储索引表访存接口 接口名称 接口含义 index_open 打开一个索引表,得到索引表的相关元信息 index_close 关闭一个索引表,释放该表的加锁或引用 index_beginscan 初始化索引扫描操作 index_beginscan_bitmap 初始化bitmap索引扫描操作 index_endscan 结束并释放索引扫描操作 index_rescan 重新开始索引扫描操作 index_markpos 记录当前索引扫描位置 index_restrpos 重置索引扫描位置 index_getnext 获取下一条符合索引条件的元组 index_getnext_tid 获取下一条符合索引条件的元组指针 index_fetch_heap 根据上面的指针,获取具体的堆表元组 index_getbitmap 获取符合索引条件的所有堆表元组指针组成的bitmap index_bulk_delete 清理索引页面上的无效元组 index_vacuum_cleanup 索引页面清理之后的统计信息和空闲空间信息更新 index_build 扫描堆表数据,构造索引表数据 和堆表存储接口不同,由于openGauss支持多种索引结构(B-Tree,hash,GIN(generalized inverted index,通用倒排索引)等),每种索引结构内部的页面间组织方式以及扫描方式都不太相同,因此在上述接口中,没有直接定义底层的页面和元组操作,而是进一步调用了各个索引自己的访存方式。不同索引的底层访存接口,可以在pg_am系统表中查询得到。 4.2.6 行存储缓存机制 行存储缓存加载和淘汰机制如图4-20所示。 图4-20 行存储缓存和淘汰机制示意图 行存储堆表和索引表页面的缓存和淘汰机制主要包含以下几个部分。 1. 共享缓冲区内存页面数组下标哈希表 共享缓冲区内存页面数组下标哈希表用于将远大于内存容量的物理页面与内存中有限个数的内存页面建立映射关系。该映射关系通过一个分段、分区的全局共享哈希表结构实现。哈希表的键值为buftag(页面标签)结构体。该结构体由“rnode”、“forkNum”、“blockNum”三个成员组成。其中“rnode”对应行存储表物理文件名的主体命名;“forkNum”对应主体命名之后的后缀命名,通过主体命名和后缀命名,可以找到唯一的物理文件;而“blockNum”对应该物理文件中的页面号。因此,该三元组可以唯一确定任意一个行存储表物理文件中的物理页面位置。哈希表的内容值为与该物理页面对应的内存页面的“buffer id”(共享内存页面数组的下标)。 因为该哈希表是所有数据页面查询的入口,所以当存在并发查询时在该哈希表上的查询和修改操作会非常频繁。为了降低读写冲突,把该哈希表进行了分区,分区个数等于NUM_BUFFER_PARTITIONS宏的定义值。在对该哈希表进行查询或修改操作之前首先需要获取相应分区的共享锁或排他锁。考虑到当对该哈希表进行插入操作时待插入的三元组键值对应的物理页面大概率不在当前的共享缓冲区中,因此该哈希表的容量等于“g_instance.attr.attr_storage.NBuffers + NUM_BUFFER_PARTITIONS”。该表具体的定义代码如下: typedef struct buftag { RelFileNode rnode; /* 表的物理文件位置结构体 */ ForkNumber forkNum; /* 表的物理文件后缀信息 */ BlockNumber blockNum; /* 页面号 */ } BufferTag; 2. 共享buffer desc数组 该数组有“g_instance.attr.attr_storage.NBuffers”个成员,与实际存储页面内容的共享buffer数组成员一一对应,用来存储相同“buffer id”(即这两个全局数组的下标)的数据页面的属性信息。该数组成员为BufferDesc结构体,具体定义代码如下: typedef struct BufferDesc { BufferTag tag; /*缓冲区页面标签 */ pg_atomic_uint32 state; /* 状态位、引用计数、使用历史计数 */ int buf_id; /*缓冲区下标 */ ThreadId wait_backend_pid; LWLock* io_in_progress_lock; LWLock* content_lock; pg_atomic_uint64 rec_lsn; volatile uint64 dirty_queue_loc; } BufferDesc; (1) tag成员是该页面的(relfilenode,forknum,blocknum)三元组。 (2) state成员是该内存状态的标志位,主要包含BM_LOCKED(该buffer desc结构体内容的排他锁标志)、BM_DIRTY(脏页标志)、BM_VALID(有效页面标志)、BM_TAG_VALID(有效tag标志)、BM_IO_IN_PROGRESS(页面I/O状态标志)等。 (3) buf_id成员,是该成员在数组中的下标。 (4) wait_backend_pid成员,是等待页面unpin(取消引用)的线程的线程号。 (5) io_in_progress_lock成员,是用于管理页面并发I/O操作(从磁盘加载和写入磁盘)的轻量级锁。 (6) content_lock成员,是用于管理页面内容并发读写操作的轻量级锁。 (7) rec_lsn成员,是上次写入磁盘之后该页面第一次修改操作的日志lsn值。 (8) dirty_queue_loc成员,是该页面在全局脏页队列数组中的(取模)下标。 3. 共享buffer数组 该数组有“g_instance.attr.attr_storage.NBuffers”个成员,每个数组成员即为保存在内存中的行存储表页面内容。需要注意的是,每个buffer在代码中以一个整型变量来标识,该值从1开始递增,数值上等于“buffer id + 1”,即“数组下标加1”。 4. bgwriter线程组 该数组有“g_instance.attr.attr_storage.bgwriter_thread_num”个线程。每个“bgwriter”线程负责一定范围内(目前为均分)的共享内存页面的写入磁盘操作,如图4-20中所示。如果全局共享buffer数组的长度为12,一共有3个“bgwriter”线程,那么第1个“bgwriter”线程负责“buffer id 0 - buffer id 3”的内存页面的维护和写入磁盘;第2个“bgwriter”线程负责“buffer id 4 - buffer id 7”的内存页面的维护和写入磁盘;第3个“bgwriter”线程负责buffer id 8 - buffer id 11的内存页面的维护和写入磁盘。每个“bgwriter”进程在后台循环扫描自己负责的那些共享内存页面和它们的buffer desc状态,将被业务修改过的脏页收集起来,批量写入双写文件,然后写入表文件系统。对于刷完的内存页,将其状态变为非脏,并追加到空闲buffer id队列的尾部,用于后续业务加载其他当前不在共享缓冲区的物理页面。每个“bgwriter”线程的信息记录在BgWriterProc结构体中,该结构体的定义代码如下: typedef struct BgWriterProc { PGPROC *proc; CkptSortItem *dirty_buf_list; uint32 dirty_list_size; int *cand_buf_list; volatile int cand_list_size; volatile int buf_id_start; pg_atomic_uint64 head; pg_atomic_uint64 tail; bool need_flush; volatile bool is_hibernating; ThrdDwCxt thrd_dw_cxt; volatile uint32 thread_last_flush; int32 next_scan_loc; } BgWriterProc; 其中比较关键的几个成员含义是: (1) dirty_buf_list为存储每批收集到的脏页面buffer id的数组。dirty_list_size为该数组的长度。 (2) cand_buf_list为存储写入磁盘之后非脏页面buffer id的队列数组(空闲buffer id数组)。cand_list_size为该数组的长度。 (3) buf_id_start为该bgwriter负责的共享内存区域的起始buffer id,该区域长度通过“g_instance.attr.attr_storage.NBuffers / g_instance.attr.attr_storage.bgwriter_thread_num”得到。 (4) head为当前空闲buffer id队列的队头数组下标,tail为当前空闲buffer id队列的队尾数组下标。 (5) next_scan_loc为上次bgwriter循环扫描时停止处的buffer id,下次收集脏页从该位置开始。 5. pagewriter线程组 “pagewriter”线程组由多个“pagewriter”线程组成,线程数量等于GUC参数(g_instance.ckpt_cxt_ctl->page_writer_procs.num)的值。“pagewriter”线程组分为主“pagewriter”线程和子“pagewriter”线程组。主“pagewriter”线程只有一个,负责从全局脏页队列数组中批量获取脏页面、将这些脏页批量写入双写文件、推进整个数据库的检查点(故障恢复点)、分发脏页给各个pagewriter线程,以及将分发给自己的那些脏页写入文件系统。子“pagewriter”线程组包括多个子“pagewriter”线程,负责将主“pagewriter”线程分发给自己的那些脏页写入文件系统。 每个“pagewriter”线程的信息保存在PageWriterProc结构体中,该结构体的定义代码如下: typedef struct PageWriterProc { PGPROC* proc; volatile uint32 start_loc; volatile uint32 end_loc; volatile bool need_flush; volatile uint32 actual_flush_num; } PageWriterProc; 其中: (1) proc成员为“pagewriter”线程属性信息。 (2) start_loc为分配给本线程待写入磁盘的脏页在全量脏页队列中的起始位置。 (3) end_loc为分配给本线程待写入磁盘的脏页在全量脏页队列中的结尾位置。 (4) need_flush为是否有脏页被分配给本“pagewriter”的标志。 (5) actual_flush_num为本批实际写入磁盘的脏页个数(有些脏页在分配给本“pagewriter”线程之后,可能被“bgwriter”线程写入磁盘,或者被DROP(删除)类操作失效)。 “pagewriter”线程与“bgwriter”线程的差别:“bgwriter”线程主要负责将脏页写入磁盘,以便留出非脏的缓冲区页面用于加载新的物理数据页;“pagewriter”线程主要的任务是推进全局脏页队列数组的进度,从而推进整个数据库的检查点和故障恢复点。数据库的检查点是数据库(故障)重启时需要回放的日志的起始位置lsn。在检查点之前的那些日志涉及的数据页面修改,需要保证在检查点推进时刻已经写入磁盘。通过推进检查点的lsn,可以减少数据库宕机重启之后需要回放的日志量,从而降低整个系统的恢复时间目标(recovery time objective,RTO)。关于“pagewriter”的具体工作原理,将在“4.2.9 持久化及故障恢复机制”小节进行更详细的描述。 6. 双写文件 一般磁盘的最小I/O单位为1个扇区(512字节),大部分文件系统的I/O单位为8个扇区。数据库最小的I/O单位为一个页面(16个扇区),因此如果在写入磁盘过程中发生宕机,可能出现一个页面只有部分数据写入磁盘的情况,会影响当前日志恢复的一致性。为了解决上述问题,openGauss引入了双写文件。所有页面在写入文件系统之前,首先要写入双写文件,并且双写文件以“O_SYNC | O_DIRECT”模式打开,保证同步写入磁盘。因为双写文件是顺序追加的,所以即使采用同步写入磁盘,也不会带来太明显的性能损耗。在数据库恢复时,首先从双写文件中将可能存在的部分写入磁盘的页面进行修复,然后再回放日志进行日志恢复。 此外也可以采用FPW(full page write,全页写)技术解决部分数据写入磁盘问题:在每次检查点之后,对于某个页面首次修改的日志中记录完整的页面数据。但是为了保证I/O性能的稳定性,目前openGauss默认使用增量检查点机制(关于增量检查点机制,参见“4.2.9 持久化及故障恢复机制”节),而该机制与FPW技术无法兼容,所以在openGauss中目前采用双写技术来解决部分数据写入磁盘问题。 结合图4-20,缓冲区页面查找的流程如下。 (1) 计算“buffer tag”对应的哈希值和分区值。 (2) 对“buffer id”哈希表加分区共享锁,并查找“buffer tag”键值是否存在。 (3) 如果“buffer tag”键值存在,确认对应的磁盘页面是否已经加载上来。如果是,则直接返回对应的“buffer id + 1”;如果不是,则尝试加载到该“buffer id”对应的缓冲区内存中,然后返回“buffer id + 1”。 (4) 如果“buffer tag”键值不存在,则寻找一个“buffer id”来进行替换。首先尝试从各个“bgwriter”线程的空闲“buffer id”队列中获取可以用来替换的“buffer id”;如果所有“bgwriter” 线程的空闲buffer id队列都为空队列,那么采用clock-sweep算法,对整个buffer缓冲区进行遍历,并且每次遍历过程中将各个缓冲区的使用计数减一,直到找到一个使用计数为0的非脏页面,就将其作为用来替换的缓冲区。 (5) 找到替换的“buffer id”之后,按照分区号从小到大的顺序,对两个“buffer tag”对应的分区同时加上排他锁,插入新“buffer tag”对应的元素,删除原来“buffer tag”对应的元素。然后再按照分区号从小到大的顺序释放上述两个分区排他锁。 (6) 最后确认对应的磁盘页面是否已经加载上来。如果是,则直接返回上述被替换的“buffer id + 1”;如果不是,则尝试加载到该“buffer id”对应的buffer内存中,然后返回“buffer id + 1”。 行存储共享缓冲区访问的主要接口和含义如表4-27所示。 表4-27 行存储共享缓冲区访问的主要接口 函数名 操作含义 ReadBufferExtended 读、写业务线程从共享缓冲区获取页面用于读、写查询 ReadBufferWithoutRelcache 恢复线程从共享缓冲区获取页面用于回放日志 ReadBufferForRemote 备机页面修复线程从共享缓冲区获取页面用于修复主机损坏页面 4.2.7 cstore 列存储格式是OLAP类数据库系统最常用的数据格式,适合复杂查询、范围统计类查询的在线分析型处理系统。本节主要介绍openGauss数据库内核中cstore列存储格式的实现方式。 1. cstore整体框架 cstore列存储格式整体框架如图4-21所示。其主要模块代码分布参见4.2.1节。与行存储格式不同,cstore列存储的主体数据文件以CU为I/O单元,只支持追加写操作,因此cstore只有读共享缓冲区。CU间和CU内的可见性由对应的CUDESE表(astore表)决定,因此其可见性和并发控制原理与行存储astore基本相同。 2. cstore存储单元结构 图4-22 CU结构示意图 如图4-22所述,cstore的存储单元是CU,分别包括以下内容。 (1) CU的CRC值,为CU结构中除CRC成员之外,其他所有字节计算出的32位CRC值。 (2) CU的magic值,为插入CU的事务号。 (3) CU的属性值,为16位标志值,包括CU是否包含NULL行、CU使用的压缩算法等CU粒度属性信息。 (4) 压缩后NULL值位图长度,如果属性值中标识该CU包含NULL行,则本CU在实际数据内容开始处包含NULL值位图,此处储存该位图的字节长度,如果该CU不包含NULL行,则无该成员。 (5) 压缩前数据长度,即CU数据内容在压缩前的字节长度,用于读取CU时进行内存申请和校验。 (6) 压缩后数据长度,即CU数据内容在压缩后的字节长度,用于插入CU时进行内存申请和校验。 (7) 压缩后NULL值位图内容,如果属性值中标识该CU包含NULL行,则该成员即为每行的NULL值位图,否则无该成员。 (8) 压缩后数据内容,即实际写入磁盘的CU主体数据内容。 每个CU最多保存对应字段的MAX_BATCH_ROWS行(默认60000行)数据。相邻CU之间按8kB对齐。 CU模块提供的主要CU操作接口如表4-28所示。 表4-28 CU操作接口 函数名称 接口含义 AppendCuData 向组装的CU中增加一行(仅对应字段) Compress 压缩(若需)和组装CU FillCompressBufHeader 填充CU头部 CompressNullBitmapIfNeed 压缩NULL值位图 CompressData 压缩CU数据 CUDataEncrypt 加密CU数据 ToVector 将CU数据解构为向量数组结构 UnCompress 解压(若需)和解析CU UnCompressHeader 解析CU头部内容 UnCompressNullBitmapIfNeed 解压NULL值位图 UnCompressData 解压CU数据 CUDataDecrypt 解密CU数据 3. cstore多版本机制 cstore支持完整事务语义的DML查询,原理如下。 (1) CU间的可见性:每个CU对应CUDESC表(astore行存储表)中的一行记录(一对一),该CU的可见性完全取决于该行记录的可见性。 (2) 同一个CU内不同行的可见性:每个CU的内部可见性对应CUDESC表中的一行(多对一),该行的bitmap字段为最长MAX_BATCH_ROWS个bit的删除位图(bit 1表示删除,bit 0表示未删除),通过该位图记录的可见性和多版本,来支持CU内不同行的可见性。同时由于DML操作都是行粒度操作的,因此对于行号范围相同的、不同字段的多个CU均对应同一行位图记录。 (3) CU文件读写并发控制:CU文件自身为APPEND-ONLY,只在追加时对文件大小扩展进行加锁互斥,无须其他并发控制机制。 (4) 同一个字段的不同CU,对应严格单调递增的cu_id编号,存储在对应的CUDESC表记录中,该cu_id的获取通过图4-24中的文件扩展锁来进行并发控制。 (5) 对于cstore表的单条插入以及更新操作,提供与每个cstore表对应的delta表(astore行存储表),来接收单条插入或单条更新的元组,以降低CU文件的碎片化。 可见,cstore表的可见性依赖于对应CUDESC表中记录的可见性。一个CUDESC表的结构如表4-29所示,其与CU的对应关系如图4-23所示。 表4-29 CUDESC表的结构 字段名 类型 含义 col_id integer 字段序号,即该cstore列存储表的第几个字段;特殊的,对于CU位图记录,该字段恒为-10 cu_id oid CU序号,即该列的第几个CU min text 该CU中该字段的最小值 max text 该CU中该字段的最大值 row_count integer 该CU中的行数 cu_mode integer CU模式 size bigint 该CU大小 cu_pointer text 该CU偏移(8k对齐);特殊的,对于CU位图记录,该字段为删除位图的二进制内容 magic integer 该CU magic号,与CU头部的magic相同,校验用 extra text 预留字段 图4-23 CUDESC表和CU对应关系示意图 如图4-24、图4-25所示,下面结合并发插入和并发插入查询2种具体场景,介绍openGauss中cstore多版本的具体实现方法。 图4-24 cstore表并发插入示意图 图4-25 cstore表并发插入和查询示意图 1)并发插入操作 对于并发的插入操作,会话1和会话2首先分别在各自的局部内存中完成待插入CU的拼接。然后假设会话1先获取到cstore表的扩展锁,那么会话2会阻塞在该锁上。在持锁阶段,会话1申请到该字段下一个cuid 1001,预占了该cu文件0 - 6 K的内容(即cuid 1001的内容大小),将cuid的大小、偏移以及cuid 1001头部部分信息填充到CUDESC记录中,并完成CUDESC记录的插入。接着,会话1放锁,并将cuid 1001的内容写入到CU对应偏移处,记录日志,再将删除位图记录插入CUDESC表中。当会话1释放cstore表的扩展锁之后,会话2就可以获取到该锁,然后,类似会话1的后续操作,完成cuid 1002的插入操作。 2) 并发插入和查询操作 假设在上述会话2的插入事务(事务号101)执行过程中,有并发的查询操作执行。对于查询操作,首先基于col_id和cuid这两个索引键对CUDESC表做索引扫描。由于事务号101在查询的快照中,因此cuid 1002的所有记录对于查询事务不可见,查询事务只能看到cuid 1001(事务号100)的那些记录。然后,查询事务根据CUDESC记录中对应的CU文件偏移和CU大小,将cuid 1001的数据从磁盘文件或缓存中加载到局部内存中,并拼接成向量数组的形式返回。 4. cstore访存接口和索引机制 cstore访存接口如表4-30所示,主要包括扫描、插入、删除和查询操作。 表4-30 cstore访存接口 接口名称 接口含义 CStoreBeginScan 开启cstore扫描 CStore::RunScan 执行cstore扫描,根据执行计划,内层执行cstore顺序扫描或者cstore min-max过滤扫描 CStoreGetNextBatch 继续扫描,返回下一批向量数组 CStoreEndScan 结束cstore扫描 CStore::CStoreScan cstore顺序扫描 CStore::CStoreMinMaxScan cstore min-max过滤扫描 CStoreInsert::BatchInsert(VectorBatch) 将输入的向量数组批量插入cstore表中 CStoreInsert::BatchInsert(bulkload_rows) 将输入的多行数组插入cstore表中 CStoreInsert::BatchInsertCommon 将一批多行数组(最多MAX_BATCH_ROWS行)插入cstore表各个列的CU文件中、插入对应CUDESC表记录、插入索引 CStoreInsert::InsertDeltaTable 将一批多行数组插入cstore表对应的delta表中 InsertIdxTableIfNeed 将一批多行数组插入cstore表的索引表中 CStoreDelete::PutDeleteBatch 将一批待删除的向量数组暂存到局部数据结构中,如果达到局部内存上限,则触发一下删除操作 CStoreDelete::PutDeleteBatchForTable CStoreDelete::PutDeleteBatch对于普通cstore表的内层实现 CStoreDelete::PutDeleteBatchForPartition CStoreDelete::PutDeleteBatch对于分区cstore表的内层实现 CStoreDelete::PutDeleteBatchForUpdate CStoreDelete::PutDeleteBatch对于更新cstore表操作的内层实现(更新操作由删除操作和插入操作组合而成) CStoreDelete::ExecDelete 执行cstore表删除,内层调用普通cstore表删除或分区cstore表删除 CStoreDelete::ExecDeleteForTable 执行普通cstore表删除 CStoreDelete::ExecDeleteForPartition 执行分区cstore表删除 CStoreDelete::ExecDelete(rowid) 删除cstore表中特定一行的接口 CStoreUpdate::ExecUpdate 执行cstore表更新 cstore表查询执行流程,可以参考图4-26中所示。其中,灰色部分实际上是在初始化cstore扫描阶段执行的,根据每个字段的具体类型,绑定不同的CU扫描和解析函数,主要有FillVector、FillVectorByTids、FillVectorLateRead3类CU扫描解析接口。 图4-26 cstore表查询流程示意图 cstore表插入执行流程,可以参考图4-27所示。其中灰色部分内的具体流程可以参考图4-24、图4-25中所示。当满足以下3个条件时,可以支持delta表插入: (1) 打开enable_delta_store GUC参数。 (2) 该批向量数组为本次导入的最后一批向量数组。 (3) 该批向量数组的行数小于delta表插入的阈值。 图4-27 cstore表插入流程示意图 cstore表的删除流程主要分为两步。 (1) 如果存在delta表,那么先从delta表中删除满足谓词条件的记录。 (2) 在CUDESC表中更新待删除行所在CU的删除位图记录。 cstore表的更新操作由删除操作和插入操作组合而成,流程不再赘述。 openGauss的cstore表支持psort和cbtree两种索引。 psort索引是一种局部排序聚簇索引。psort索引表的组织形式也是cstore表,该cstore表的字段包括索引键中的各个字段,再加上对应的行号(TID)字段。如图4-28所示,将一定数量的记录按索引键经过排序聚簇之后,与TID字段共同拼装成向量数组之后,插入psort索引cstore表中,插入流程和上面cstore表插入流程相同。 图4-28 psort索引插入原理图 查询时如果使用psort索引扫描,会首先扫描psort索引cstore表(扫描方式和上面cstore表扫描流程相同)。在一个psort索引CU的内部,由于做了局部聚簇索引,因此可以使用基于索引键的二分查找方式,快速找到符合索引条件的记录在该psort索引中的行号,该行的TID字段值即为该条记录在cstore主表中的行号。上述流程如图4-29所示。值得一提的是由于做了局部聚簇索引,因此在索引cstore表扫描过程中,在真正加载索引表CU文件之前,可以通过CUDESC中的min max做到非常高效的初筛过滤。 cstore表的cbtree索引和行存储表的B-Tree索引在结构和使用方式上几乎完全一致,相关原理可以参考行存储索引章节(“4.2.5 行存储索引机制”节),此处不再赘述。 openGauss cstore表索引对外提供的主要接口如表4-31所示。 表4-31 cstore表索引对外接口 接口名称 接口含义 psortgettuple 通过psort索引,返回下一条满足索引条件的元组。伪接口,实际psort索引扫描通过CStore::RunScan实现 psortgetbitmap 通过psort索引,返回满足索引条件的元组的tid bitmap。伪接口,实际psort索引扫描通过CStore::RunScan实现 psortbuild 构建psort索引表数据。主要流程包括,从cstore主表中扫描数据、局部聚簇排序、插入到psort索引cstore表中 cbtreegettuple 通过cbtree索引,返回下一条满足索引条件的元组。内部和btgettuple都是通过调用_bt_gettuple_internal函数实现的 cbtreegetbitmap 通过cbtree索引,返回满足索引条件的元组的tid bitmap。内部和btgetbitmap都是通过调用_bt_next函数实现的 cbtreebuild 构建cbtree索引表数据。内部实现与btbuild类似,先后调用_bt_spoolinit、CStoreGetNextBatch、_bt_spool、_bt_leafbuild和_bt_spooldestroy等几个主要函数实现。与btbuild区别在于,B-Tree的构建过程中,扫描堆表是通过heapam接口实现的,而cbtree扫描的是cstore表,因此使用的是CStoreGetNextBatch 5. cstore缓存机制 考虑到cstore列存储格式主要面向只读查询居多的OLAP类业务,因此openGauss提供只读的共享CU缓冲区机制。 openGauss中CU只读共享缓冲区的结构如图4-30所示。和行存储页面粒度的共享缓冲区类似,最上层为共享哈希表,哈希表键值为CU的slot类型、relfilenode、colid、cuid、cupointer构成的五元组,哈希表的记录值为该CU对应的缓冲区槽位slot id(对应行存储共享缓区的buffer id)。在全局CacheDesc数组中,用CacheDesc结构体记录与slot id对应的缓存槽位的状态信息(对应行存储缓冲区的BufferDesc结构体)。在共享CU数组中,用CU结构体记录与slot id对应的缓存CU的结构体信息。 与行存储固定的页面大小不同,不同CU的大小可能是不同的(行存储页面大小都是8 K),因此上述CU槽位只记录指向实际内存中CU数据的指针。另一方面为了保证共享内存大小可控,通过另外的全局变量来记录已经申请的有效槽位中所有CU的大小总和。 图4-30 CU只读共享缓存结构示意图 CU只读共享缓冲区的工作机制如图4-31所示。 (1) 当从磁盘读取一个CU放如Cache Mgr时,需要从FreeSlotList里拿到一个free slot(空闲槽位)存放CU,然后插入到哈希表中。 (2) 当FreeSlotList为NULL的时,需要根据LRU算法淘汰掉一个slot(槽位),释放CU data占的内存,减小CU总大小计数,并从哈希表中删除,然后存放新的CU,再插入哈希表中。 (3) 缓存大小可以配置。如果内存超过设置的大小,需要淘汰掉适量的slot,并释放CU data占用的内存。 (4) 支持缓存压缩态的CU或解压态的CU两种模式,可以通过配置文件修改,同时只能存在一种模式。 图4-31 CU只读共享缓存读取示意图 与CU只读共享缓冲区相关的关键数据结构代码如下: typedef struct CUSlotTag { RelFileNodeOld m_rnode; int m_colId; int32 m_CUId; uint32 m_padding; CUPointer m_cuPtr; } CUSlotTag; /* slot id哈希表键值主要部分,各个成员的含义从命名中可以清晰看出 */ typedef struct DataSlotTag { DataSlotTagKey slotTag; CacheType slotType; } DataSlotTag; /* slot id哈希表键值结构体,成员包括CUSlotTag与slot类型(CU、OBS外表等) */ typedef struct CacheLookupEnt { CacheTag cache_tag; CacheSlotId_t slot_id; } CacheLookupEnt; /* slot id哈希表记录结构体,成员包括哈希表键值和对应的slot id */ typedef struct CacheDesc { uint16 m_usage_count; uint16 m_ring_count; uint32 m_refcount; CacheTag m_cache_tag; CacheSlotId_t m_slot_id; CacheSlotId_t m_freeNext; LWLock *m_iobusy_lock; LWLock *m_compress_lock; /*The data size in the one slot.*/ int m_datablock_size; bool m_refreshing; slock_t m_slot_hdr_lock; CacheFlags m_flag; } CacheDesc; /* CU共享缓冲区槽位状态结构体,其中m_usage_count、m_ring_count为LRU淘汰算法需要的使用计数,m_refcount为判断能否淘汰的被引用计数,m_freeNext指向下一次空闲的slot槽位(如果本槽位在free list中的话,否则m_freeNext恒等于-2),m_iobusy_lock为I/O并发控制锁,m_compress_lock为压缩并发控制锁,m_datablock_size为CU实际数据的大小,m_slot_hdr_lock保护整个CacheDesc的并发读写操作,m_flag表示槽位状态(包括全新、有效、freelist中、空闲、I/O中、错误等状态)*/ 4.2.8 日志系统 内存是一种易失性存储介质,在断电等场景下存储在内存介质中的数据会丢失。为了保障数据的可靠性需要将共享缓冲区中的脏页写入磁盘,此即数据的持久化过程。对于最常用的持久化存储介质磁盘,由于每次读写操作都有一个“启动”代价,导致磁盘的读写操作频率有一个上限。即使是超高性能的SSD磁盘,其读写频率也只能达到10000次/秒左右。如果多个磁盘读写请求的数据在磁盘上是相邻的,就可以被合并为一次读写操作。因为合并后可以等效降低读写频率,所以磁盘顺序读写的性能通常要远优于随机读写。由于如上原因,数据库通常都采用顺序追加的预写日志(write ahead log,WAL)来记录用户事务对数据库页面的修改。对于物理表文件所对应的共享内存中的脏页会等待合适的时机再异步、批量地写入磁盘。 日志可以按照用户对数据库不同的操作类型分为以下几类,每种类型日志分别对应一种资源管理器,负责封装该日志的子类、具体结构以及回放逻辑等。如表4-32所示。 表4-32 日志类型 日志类型名字 资源管理器类型 对应操作 XLOG RM_XLOG_ID pg_control控制文件修改相关的日志,包括检查点推进、事务号分发、参数修改、备份结束等 Transaction RM_XACT_ID 事务控制类日志,包括事务提交、回滚、准备、提交准备、回滚准备等 Storage RM_SMGR_ID 底层物理文件操作类日志,包括文件的创建和截断 CLOG RM_CLOG_ID 事务日志修改类日志,包括CLOG拓展、CLOG标记等 Database RM_DBASE_ID 数据库DDL类日志,包括创建、删除、更改数据库等 Tablespace RM_TBLSPC_ID 表空间DDL类日志,包括创建、删除、更新表空间等 MultiXact RM_MULTIXACT_ID MultiXact类日志,包括MultiXact槽位的创建、成员页面的清空、偏移页面的清空等 RelMap RM_RELMAP_ID 表文件名字典文件修改日志 Standby RM_STANDBY_ID 备机支持只读相关日志 Heap RM_HEAP_ID 行存储文件修改类日志,包括插入、删除、更新、pd_base_xid修改、新页面、加锁等操作 Heap2 RM_HEAP2_ID 行存储文件修改类日志,包括空闲空间清理、元组冻结、元组可见性修改、批量插入等 Heap3 RM_HEAP3_ID 行存储文件修改类日志,目前该类日志不再使用,后续可以拓展 Btree RM_BTREE_ID B-Tree索引修改相关日志,包括插入、节点分裂、插入叶子节点、空闲空间清理等 hash RM_HASH_ID hash索引修改相关日志 Gin RM_GIN_ID GIN索引(generalized inverted index,通用倒排索引)修改相关日志 Gist RM_GIST_ID Gist索引修改相关日志 SPGist RM_SPGIST_ID SPGist索引相关日志 Sequence RM_SEQ_ID 序列修改相关日志,包括序列推进、属性更新等 Slot RM_SLOT_ID 流复制槽修改相关日志,包括流复制槽的创建、删除、推进等 MOT RM_MOT_ID 内存引擎相关日志 openGauss日志文件、页面和日志记录的格式如图4-32所示。 图4-32 日志文件、页面和记录格式示意图 日志文件在逻辑意义上是一个最大长度为64位无符号整数的连续文件。在物理分布上,该逻辑文件按XLOG_SEG_SIZE大小(默认为16MB)切断,每段日志文件的命名规则为“时间线+日志id号+该id内段号”。“时间线”用于表示该日志文件属于数据库的哪个“生命历程”,在时间点恢复功能中使用。“日志id号”从0开始,按每4G大小递增加1。“id内段号”表示该16MB大小的段文件在该4G“日志id号”内是第几段,范围为0至255。上面3个值在日志段文件名中都以16进制方式显示。 每个日志段文件都可以用XLOG_BLCKSZ(默认8kB)为单位,划分为多个页面。每个8kB页面中,起始位置为页面头,如果该页是整个段文件的第一个页面,那么页面头为一个长页头(XLogLongPageHeader),否则为一个正常页头(短页头)(XLogPageHeader)。在页头之后跟着一条或多条日志记录。每个日志记录对应一个数据库的某种操作。为了降低日志记录的大小(日志写入磁盘时延是影响事务时延的主要因素之一),每条日志内部都是紧密排列的。各条日志之间按8字节(64位系统)对齐。一条日志记录可以跨两个及以上的日志页面,其最大长度限制为1G。对于跨页的日志记录,其后续日志页面页头的标志位XLP_FIRST_IS_CONTRECORD会被置为1。 长、短页头结构体的定义如下,其中存储了用于校验的magic信息、页面标志位信息、时间线信息、页面(在整个逻辑日志文件中的)偏移信息、有效长度信息、系统识别号信息、段尺寸信息、页尺寸信息等。 短页头结构体的代码如下: typedef struct XLogPageHeaderData { uint16 xlp_magic; /* 日志magic校验信息 */ uint16 xlp_info; /* 标志位 */ TimeLineID xlp_tli; /* 该页面第一条日志的时间线 */ XLogRecPtr xlp_pageaddr; /* 该页面起始位置的lsn */ uint32 xlp_rem_len; /*如果是跨页记录,本字段描述该跨页记录在本页面内的剩余长度 */ } XLogPageHeaderData; 长页头结构体的代码如下: typedef struct XLogLongPageHeaderData { XLogPageHeaderData std; /* 短页头 */ uint64 xlp_sysid; /* 系统标识符,和pg_control文件中相同 */ uint32 xlp_seg_size; /* 单个日志文件的大小 */ uint32 xlp_xlog_blcksz; /* 单个日志页面的大小 */ } XLogLongPageHeaderData; 单条日志记录的结构如图4-32中所示,其由5个部分组成: (1) 日志记录头,对应XLogRecord结构体,存储了记录长度、主备任期号、事务号、上一条日志记录起始偏移、标志位、所属的资源管理器、crc校验值等信息。 (2) 1 - 33个相关页面的元信息,对应XLogRecordBlockHeader结构体,存储了页面下标(0 - 32)、页面对应的物理文件的后缀、标志位、页面数据长度等信息;如果该日志没有对应的页面信息,则无该部分。 (3) 日志数据主体的元信息,对应(长/短)XLogRecordDataHeader结构体,记录了特殊的页面下标,用于和第二部分区分,以及主体数据的长度。 (4) 1 - 33个相关页面的数据;如果该日志没有对应的页面信息,则无该部分。 (5) 日志数据主体。 这5部分对应的结构体代码如下。如上所述,在记录日志内容时,每个部分之间是紧密挨着的,无补空字符。如果一个日志记录没有对应的相关页面信息,那么第2和第4部分将被跳过。 typedef struct XLogRecord { uint32 xl_tot_len; /* 记录总长度 */ uint32 xl_term; TransactionId xl_xid; /* 事务号 */ XLogRecPtr xl_prev; /* 前一条记录的起始位置lsn */ uint8 xl_info; /* 标志位 */ RmgrId xl_rmid; /* 资源管理器编号 */ int2 xl_bucket_id; pg_crc32c xl_crc; /* 该记录的CRC校验值 */ /* 后面紧接XLogRecordBlockHeaders或XLogRecordDataHeader结构体 */ } XLogRecord; typedef struct XLogRecordBlockHeader { uint8 id; /* 页面下标(即该记录中包含的第几个页面信息) */ uint8 fork_flags; /* 页面属于哪个后缀文件,以及标志位 */ uint16 data_length; /* 实际页面相关的数据长度(紧接该头部结构体) */ /* 如果BKPBLOCK_HAS_IMAGE标志位为1,后面紧跟XLogRecordBlockImageHeader结构体以及页面内连续数据 */ /* 如果BKPBLOCK_SAME_REL标志位没有设置,后面紧跟RelFileNode结构体 */ /* 后面紧跟页面号 */ } XLogRecordBlockHeader; typedef struct XLogRecordDataHeaderShort { uint8 id; /* 特殊的XLR_BLOCK_ID_DATA_SHORT页面下标 */ uint8 data_length; /* 短记录数据长度 */ } XLogRecordDataHeaderShort; typedef struct XLogRecordDataHeaderLong { uint8 id; /* 特殊的XLR_BLOCK_ID_DATA_LONG页面下标 */ /* 后面紧跟长记录长度,无对齐 */ } XLogRecordDataHeaderLong; 单条日志记录的操作接口主要分为插入(写)和读接口。其中,一个完整的日志插入操作一般包含以下几步接口,如表4-33所示。 表4-33 日志插入操作 步骤序号 接口名称 对应操作 1 XLogBeginInsert 初始化日志插入相关的全局变量 2 XLogRegisterData 注册该日志记录的主体数据 3 XLogRegisterBuffer/ XLogRegisterBlock 注册该日志记录相关页面的元信息 4 XLogRegisterBufData 注册该日志记录相关页面的数据 5 XLogInsert 执行真正的日志插入,包含5.1和5.2 5.1 XLogRecordAssemble 将上述注册的所有日志信息,按照图4-32中所示的紧密排列的5部分,重新组合成完整的二进制串 5.2 XLogInsertRecord 在整个逻辑日志中,预占偏移和长度,计算CRC,将完整的日志记录拷贝到日志共享缓冲区中 日志的读接口为XLogReadRecord接口。该接口从指定的日志偏移处(或上次读到的那条记录结尾位置处)开始读取和解析下一条完整的日志记录。如果当前缓存的日志段文件页面中无法读完,那么会调用ReadPageInternal接口加载下一个日志段文件页面到内存中继续读取,直到读完所有等于日志头部xl_tot_len长度的日志数据。然后,调用DecodeXLogRecord接口,将日志记录按图4-32中所示的5个组成部分进行解析。 日志文件读写的最小I/O粒度为一个页面。在事务执行过程中,只会进行(顺序追加)写日志操作。为了提高写日志的性能,在共享内存中,单独开辟一片特定大小的区域,作为写日志页面的共享缓冲区。对该共享缓冲区的并发操作(拷贝日志记录到单个页面中、淘汰lsn过老的页面、读取单个页面并写入磁盘)是事务执行流程中的关键瓶颈之一,对整个数据库系统的并发能力至关重要。 图4-33 并发日志写入流程示意图 如图4-33所示,在openGauss中对该共享缓冲区的操作采用Numa-aware的同步机制,具体步骤如下。 (1) 业务线程在本地内存中将日志记录组装成图4-32中所示的、5部分组成的字节流 (2) 找到本线程所绑定的NUMA Node对应的日志插入锁组,并在该锁组中随机找一个槽位对应的锁。 (3) 检查该锁的组头线程号。如果没有说明本线程是第一个请求该锁的,那么这个锁上所有的写日志请求将由本线程来执行,将锁的组头线程号设置为本线程号;否则说明已经存在这批写日志请求的组头线程,记录下当前组头线程的线程号,并将自己加入到这批的插入组队列中,等待组头线程完成日志插入。 (4) 对于组头线程,获取该日志插入锁的排他锁。 (5) 为该组所有的插入线程在逻辑日志文件中占位,即对当前该文件的插入偏移进行原子CAS(compare and swap,比较后交换)操作。 (6) 将该组所有后台线程本地内存中的日志依次拷贝到日志共享缓冲区的对应页面中。每当需要拷贝到下一个共享内存页面时,需要判断下一个页面对应的逻辑页面号是否和插入者的预期页面号一致(因为共享内存有限,因此同一个共享内存页面对应取模相同的逻辑页面)。首先,将自己预期的逻辑页面号,写入当前持有的日志插入锁的槽位中,然后进行上述判断。如果不一致,即共享内存页面当前的逻辑页面号比插入者预期的逻辑页面号要小,那么需要将该页面数据从共享内存中写入到磁盘,然后才能复用为新的逻辑页面号。为了防止可能还有并发业务线程在向该共享内存页面拷贝属于当前逻辑页面号的日志数据,因此需要阻塞遍历每个日志插入者持有的插入锁,直到日志插入锁被释放,或者被持有的插入锁的逻辑页面号大于目标共享内存页面中现有的逻辑页面号。经过上述检查之后,就可以保证没有并发的业务线程还在对该共享内存页面写入对应当前逻辑页面号的日志数据,因此可以将其内容写入磁盘,并更新其对应的逻辑页面号为目标逻辑页面号。 (7) 重复上一步操作,直到把该组所有后台线程待插入的日志记录拷贝完。 (8) 释放日志插入锁。 (9) 唤醒本组所有后台线程。 4.2.9 持久化及故障恢复机制 1. 行存储持久化和检查点机制 如“4.2.8 日志系统”节中所述,通过采用WAL日志的方式可以在对性能影响较小的情况下保障用户事务对数据库修改的持久化。然而如果只是依赖日志来保障持久化的话,那么数据库服务(故障)重启之后将需要回放大量的日志数据量,这会导致很大的RTO,对业务的可用性影响极大。因此共享缓冲区中的脏页也需要异步地写入磁盘中,来减少宕机重启后所需要回放的日志数据量,降低系统的RTO时间。 如果数据库系统在事务提交之后、异步写入磁盘的脏页写入磁盘之前发生宕机,那么需要在数据库再次启动之后,首先把那些宕机之前还没有来得及写入磁盘的脏页上的修改所对应的日志进行回放,使得这些脏页可以恢复到宕机之前的内容。 基于如上原理,可以得出数据库持久化的一个关键是:在宕机重启的时候,通过某种机制确定从WAL的哪个lsn开始进行恢复;可以保证在该lsn之前的那些日志,它们涉及的数据页面修改已经在宕机之前完成写入磁盘。这个恢复起始的lsn,即是数据库的检查点。 在“4.2.6 行存储缓存机制”节介绍行存储缓存加载和淘汰机制中,已经知道参与脏页写入磁盘的主要有两类线程:bgwriter和pagewriter。前者负责脏页持久化的主体工作;后者负责数据库检查点lsn的推进。openGauss采用一个无锁的全局脏页队列数组来依次记录曾经被用户写操作置脏的那些数据页面。该队列数组成员为DirtyPageQueueSlot结构体,定义代码如下,其中:buffer为队列成员对应的buffer(该值为buffer id + 1),slot_state为该队列成员的状态。 typedef struct DirtyPageQueueSlot { volatile int buffer; pg_atomic_uint32 slot_state; } DirtyPageQueueSlot; 图4-34 全局脏页队列的运行机制和检查点的推进机制 全局脏页队列的运作机制如图4-34所示,它的实现方式是一个多生产者、单消费者的循环数组。单个/多个业务线程是脏页队列的生产者,在其要修改数据页面之前,首先判断该页面buffer desc的首次脏页lsn是否非0:若该脏页buffer desc中的首次脏页lsn已经非0,说明该脏页在之前置脏的时候就已经被加入到脏页队列中,那么本次就跳过加入脏页队列的步骤;否则,对当前脏页队列的tail位置进行CAS加1操作,完成队列占位,同时,在上述CAS操作中,获取了脏页队列的lsn位置lsn1。然后,将占据的槽位位置(即CAS之前的tail值)和lsn1记录到脏页的buffer desc中。接着,将脏页的“buffer id”记录到占位的槽位中,再将槽位状态置为valid。最后,记录页面修改的日志,并尝试将该日志的位置lsn2更新到脏页队列的lsn中(如果此时脏页队列的lsn值已经被其他写业务更新为更大的值,则本线程就不更新了,也是一个CAS操作)。 基于上面这种机制,当将脏页队列中某个成员对应的脏页写入磁盘之后,检查点即可更新到该脏页“buffer desc”中记录的lsn位置。小于该lsn位置的日志,它们对应修改的页面,已经在记录这些日志之前就被加入到脏页队列中,亦即这些脏页在全局脏页队列中的位置一定比当前脏页更靠前,因此一定已经保证写入磁盘了。在图4-34中,“pagewriter”线程作为全局脏页队列唯一的消费者,负责从脏页队列中批量获取待写入磁盘的脏页,在完成写入磁盘操作之后,“pagewriter”自身不负责检查点的推进,而只是推进整个脏页队列的队头到下一个待写入磁盘的槽位位置。 实际检查点的推进由“Checkpointer”线程来负责。这是因为“pagewriter” 线程的写入磁盘操作,只是将共享缓冲区中的脏页写入到文件系统的缓存中,(由于文件系统的I/O合并优化)此时可能并没有真正写入磁盘。因此,在“Checkpointer”线程中,其先获取当前全局脏页队列的队头位置,以及对应槽位中脏页的首次脏页lsn值,然后对截至目前所有被写入文件系统的文件进行fsync(刷盘)操作,保证文件系统将它们写入物理磁盘中。然后就可以将上述lsn值作为检查点位置更新到control文件中,用于数据库重启之后回放日志的起始位置。 上述这套持久化和检查点推进机制的主要控制信息,保存在knl_g_ckpt_context结构体中,该结构体定义代码如下: typedef struct knl_g_ckpt_context { uint64 dirty_page_queue_reclsn; uint64 dirty_page_queue_tail; CkptSortItem* CkptBufferIds; /* 脏页队列相关成员 */ DirtyPageQueueSlot* dirty_page_queue; uint64 dirty_page_queue_size; pg_atomic_uint64 dirty_page_queue_head; pg_atomic_uint32 actual_dirty_page_num; /* pagewriter线程相关成员 */ PageWriterProcs page_writer_procs; uint64 page_writer_actual_flush; volatile uint64 page_writer_last_flush; /* 全量检查点相关信息成员 */ volatile bool flush_all_dirty_page; volatile uint64 full_ckpt_expected_flush_loc; volatile uint64 full_ckpt_redo_ptr; volatile uint32 current_page_writer_count; volatile XLogRecPtr page_writer_xlog_flush_loc; volatile LWLock *backend_wait_lock; volatile bool page_writer_can_exit; volatile bool ckpt_need_fast_flush; /* 检查点刷页相关统计信息(除数据页面外) */ int64 ckpt_clog_flush_num; int64 ckpt_csnlog_flush_num; int64 ckpt_multixact_flush_num; int64 ckpt_predicate_flush_num; int64 ckpt_twophase_flush_num; volatile XLogRecPtr ckpt_current_redo_point; uint64 pad[TWO_UINT64_SLOT]; } knl_g_ckpt_context; 其中和当前上述检查点机制相关的成员有: (1) dirty_page_queue_reclsn是脏页队列的lsn位置,dirty_page_queue_tail是脏页队列的队尾,这两个成员构成一个16字节的整体,通过128位的CAS操作进行整体原子读、写操作,保证脏页队列中每个成员记录的lsn一定随着入队顺序单调递增。 (2) CkptBufferIds是每批pagewriter待刷脏页数组。 (3) dirty_page_queue是全局脏页队列数组。 (4) dirty_page_queue_size是脏页数组长度,等于“g_instance.attr.attr_storage.NBuffers * PAGE_QUEUE_SLOT_MULTI_NBUFFERS”,当前PAGE_QUEUE_SLOT_MULTI_NBUFFERS取值5,以防止脏页队列因为DDL(data definition language,数据定义语言)等操作引入的空洞过多,导致脏页队列撑满阻塞业务的场景。 (5) dirty_page_queue_head是脏页队列头部。 (6) actual_dirty_page_num是脏页队列中实际的脏页数量。 2. 故障恢复机制 当数据库发生宕机重启之后需要从检查点位置开始回放之后所有的日志。不同类型的日志的回放逻辑由对应的资源管理器来实现。 当用户业务压力较大时会同时有很多业务线程并发执行事务和日志记录的插入,单位时间内产生的日志量是非常大的。对此openGauss采用多种回放线程组来进行日志的并行回放,各个回放线程组之间采用高效的流水线工作方式,各个回放线程组内采用多线程并行的工作方式,以便保证日志的回放速率不会明显低于日志产生的速率。 图4-35 openGauss并行回放流程示意图 openGauss并行回放流程如图4-35所示,其中每个线程(组)的运行机制如下。 (1) “Walreceiver”线程收到日志成功写入磁盘后,“XLogReadWorker”线程从“Walreceiver”线程的缓冲区中读取字节流,“XLogReadManager”线程将字节流decode(解码)成redoitem(单个回放对象)。“Startupxlog”线程按照表文件名粒度(refilenode)将redoitem发放给各个“ParseRedoRecord”线程,其他的日志发送给“TrxnManager”线程。 (2) “ParseRedoRecord”线程负责表文件(relation)相关的日志处理,从队列中获取批量的日志进行解析,将日志按照页面粒度进行拆分,然后发给“PageRedoManager”线程。拆分原理如下。 针对行存储表、索引等数据页面操作的日志,按照涉及的页面个数拆成多条日志。例如heap_update日志,如果删除的老元组和插入的新元组在不同的页面上,那么会被拆成2条,分别插入到哈希表中。 xact、truncate、drop database等日志是针对表的,不能进行拆分。针对这些日志,先清理掉哈希表中相关日志,然后等这些日志之前的日志都回放之后,再在PageRedoManger中进行回放,并将该日志分发给所有“PageRedoWorker”线程来进行invalid page(无效页面)的清理、数据写入磁盘等操作。 针对Createdb(创建数据库)操作要等所有“PageRedoWorker”线程将Createdb日志之前的日志都回放后,再由一个“PageRedoManager”线程进行Createdb操作的回放。这个过程中其余线程需要等待Createdb操作回放结束后才能继续回放后续日志。 (3) “PageRedoManager”线程利用哈希表按照页面粒度组织日志,同一个页面的日志按照lsn顺序放入到一个列表中,之后将页面日志列表分发给“PageRedoWorker”线程。 (4) “PageRedoWorker”线程负责页面日志回放功能,从队列中获取一个日志列表进行批量处理。 (5) “TrxnManager”线程负责事务相关的XLOG日志的分发,以及需要全局协调的事务处理。 (6) “TrxnWorker”线程负责事务日志回放功能,从队列中获取一个日志进行处理。当前只有一个“TrxnWorker”线程负责处理事务日志。 为了保证高效的日志分发性能,“PageRedoManager”进程和“PageRedoWorker”进程之间采用了带阻塞功能的无锁单生产者单消费者(single producer single consumer,SPSC)队列。如图4-36所示,分配线程作为生产者将解析后的日志放入回放线程的列队中,回放线程从队列中消费日志进行回放。另一方面为了提升整体并行回放机制的可靠性,会在一个页面的回放动作中对日志记录头部的lsn和页面头部的lsn进行校验,以保证回放过程中数据库系统的一致性。 图4-36 无锁SPSC队列示意图 3. cstore列存储持久化机制 由于在openGauss中cstore主体数据没有写缓冲区,因此对于所有的插入或更新事务,在拼装完新的CU之后都是直接调用pwrite来写入文件系统缓存,并且在事务提交之前调用“CUStorage::FlushDataFile”接口完成本地磁盘的持久化(该函数内部调用fsync执行写入磁盘)。由于OLAP系统中通常插入事务都是批量导入执行的,因此在这个过程中对于cstore表物理文件的写操作基本都是顺序I/O,可以获得较高的性能。 4.2.10 主备机制 openGauss提供主备机制来保障数据的高可靠和数据库服务的高可用。如图4-37所示,在主、备实例之间通过日志复制来进行数据库数据和状态的一致性同步。日志同步是指将主机对数据的修改日志同步到备机,备机通过日志回放将日志重新还原为数据修改。 图4-37 主备机日志同步示意图 参与日志同步的主要有“wal sender”(主机端)和“wal receiver”(备机端)两个线程。一个主机上可以由多个“wal sender”线程同时存在,用于给不同的备机进行日志复制;一个备机上同一时刻只会有一个“wal receiver”线程,从唯一一个指定的主机上拷贝日志。 “wal sender”线程的所有关键信息均保存在knl_t_walsender_context结构体中,其定义代码如下: typedef struct knl_t_walsender_context { char* load_cu_buffer; int load_cu_buffer_size; struct WalSndCtlData* WalSndCtl; struct WalSnd* MyWalSnd; int logical_xlog_advanced_timeout; DemoteMode Demotion; bool wake_wal_senders; bool wal_send_completed; int sendFile; XLogSegNo sendSegNo; uint32 sendOff; struct WSXLogJustSendRegion* wsXLogJustSendRegion; XLogRecPtr sentPtr; XLogRecPtr catchup_threshold; struct StringInfoData* reply_message; struct StringInfoData* tmpbuf; char* output_xlog_message; Size output_xlog_msg_prefix_len; char* output_data_message; uint32 output_data_msg_cur_len; XLogRecPtr output_data_msg_start_xlog; XLogRecPtr output_data_msg_end_xlog; struct XLogReaderState* ws_xlog_reader; TimestampTz last_reply_timestamp; TimestampTz last_logical_xlog_advanced_timestamp; bool waiting_for_ping_response; volatile sig_atomic_t got_SIGHUP; volatile sig_atomic_t walsender_shutdown_requested; volatile sig_atomic_t walsender_ready_to_stop; volatile sig_atomic_t response_switchover_requested; ServerMode server_run_mode; char gucconf_file[MAXPGPATH]; char gucconf_lock_file[MAXPGPATH]; FILE* ws_dummy_data_read_file_fd; uint32 ws_dummy_data_read_file_num; struct cbmarray* CheckCUArray; struct LogicalDecodingContext* logical_decoding_ctx; XLogRecPtr logical_startptr; int remotePort; bool walSndCaughtUp; } knl_t_walsender_context; 其中: (1) WalSndCtl指向保存全局所有“wal sender”线程控制状态的共享结构体,是一致性复制协议的关键所在。 (2) MyWalSnd指向上述全局共享结构体中当前“wal sender”线程的槽位。 (3) Demotion为当前主机降备模式,分为未降备(NoDemote)、优雅降备(SmartDemote)和快速降备(FastDemote)。 (4) sendFile、sendSegNo、sendOff用于保存当前复制的日志文件的文件操作状态。 (5) reply_message用于保存备机回复的消息。 (6) output_xlog_message为待发送的日志内容主体。 (7) server_run_mode为wal sender线程启动时的HA(high availability,HA)高可靠性)状态,即主机(primary)、备机(standby)或未决(pending)。 (8) walSndCaughtUp指示备机是否已经追赶上主机。 (9) remotePort为wal receiver线程的端口,用于身份验证。 (10) load_cu_buffer 、load_cu_buffer_size 、output_data_message、output_data_msg_cur_len、output_data_msg_start_xlog、output_data_msg_end_xlog、ws_xlog_reader、CheckCUArray为后续支持混合类型(日志+增量页面)复制的预留接口。 wal receiver线程的所有关键信息均保存在knl_t_walreceiver_context结构体中,其定义代码如下: typedef struct knl_t_walreceiver_context { volatile sig_atomic_t got_SIGHUP; volatile sig_atomic_t got_SIGTERM; volatile sig_atomic_t start_switchover; char gucconf_file[MAXPGPATH]; char temp_guc_conf_file[MAXPGPATH]; char gucconf_lock_file[MAXPGPATH]; char** reserve_item; time_t standby_config_modify_time; time_t Primary_config_modify_time; TimestampTz last_sendfilereply_timestamp; int check_file_timeout; struct WalRcvCtlBlock* walRcvCtlBlock; struct StandbyReplyMessage* reply_message; struct StandbyHSFeedbackMessage* feedback_message; struct StandbySwitchRequestMessage* request_message; struct ConfigModifyTimeMessage* reply_modify_message; volatile bool WalRcvImmediateInterruptOK; bool AmWalReceiverForFailover; bool AmWalReceiverForStandby; int control_file_writed; } knl_t_walreceiver_context; 其中: (1) walRcvCtlBlock指向“wal receiver”线程主控数据,保存当前日志复制进度,备机日志写盘、写入磁盘进度,以及接收日志缓冲区。 (2) reply_message保存用于回复主机的消息。 (3) feedback_message用于保存热备的相关信息,供主机空闲空间清理时参考。 (4) request_message用于保存主机降备请求的相关信息。 (5) reply_modify_message用于保存请求配置文件复制的相关信息。 (6) AmWalReceiverForFailover表示当前“wal receiver”线程处于failover场景下,连接从备进行日志追赶。 (7) AmWalReceiverForStandby表示当前“wal receiver”线程为连接备机进行日志复制的级联备机。 主备日志同步,主要包括以下6个场景。 1. 备机发起复制请求,进入流式复制。 图4-38 主备建连和流式复制流程图 如图4-38所示,日志复制请求是由“wal receiver”线程发起的。在libpqrcv_connect函数中,备机通过libpq协议连上主机,通过特殊的连接串信息,触发主机侧启动“wal sender”线程来处理该连接请求(相比之下,对于普通客户端查询请求,主机启动backend线程或线程池线程来处理连接请求)。在WalSndHandshake函数中,wal sender线程与wal receiver线程完成身份、日志一致性等校验之后,进入WalSndLoop开始日志复制循环。主要的主、备机握手和校验报文如表4-34所示,在主机收到T_StartReplicationCmd报文之后,开始进入日志复制阶段。 表4-34 主、备机握手和校验报文 报文类型 报文作用 T_IdentifySystemCmd 请求主机发送主机侧system_identifier,校验是否和备机一致 T_IdentifyVersionCmd 请求主机发送主机侧版本号,校验是否和备机一致 T_IdentifyModeCmd 请求主机发送主机侧HA状态,校验是否是主机状态 T_IdentifyMaxLsnCmd 请求主机发送当前最大的lsn位置(即日志偏移),用于备机重建 T_IdentifyConsistenceCmd 请求主机发送指定lsn位置日志记录的crc值,校验是否和备机一致 T_IdentifyChannelCmd 请求主机校验备机的端口是否在repliconn_info参数中,返回校验结果 T_IdentifyAZCmd 请求主机发送主机侧AZ名字 T_BaseBackupCmd 请求主机开始发起全量重建 T_CreateReplicationSlotCmd 请求主机创建流复制槽 T_DropReplicationSlotCmd 请求主机删除流复制槽 T_StartReplicationCmd 请求主机开始日志复制 2. Quorum一致性复制协议 为了保证数据库数据的可靠和高可用,当主机上执行的事务修改产生日志之后,在事务提交之前需要将本事务产生的日志同步到多个备机上。openGauss采用Quorum一致性复制协议,即当多数备机完成上述事务的日志同步之后主机事务方可提交。这个过程中作为事务提交参考的是同步备,其他备机是异步备,作为冗余备份。同步备和异步备的具体选择可以通过配置synchronus_standby_names参数实现。 图4-39 事务提交和一致性复制协议 主机上事务提交和一致性复制协议的工作运行机制如图4-39所示。主要涉及的数据结构是WalSndCtlData数据结构体,其定义代码如下: typedef struct WalSndCtlData { SHM_QUEUE SyncRepQueue[NUM_SYNC_REP_WAIT_MODE]; XLogRecPtr lsn[NUM_SYNC_REP_WAIT_MODE]; bool sync_standbys_defined; bool most_available_sync; bool sync_master_standalone; DemoteMode demotion; slock_t mutex; WalSnd walsnds[FLEXIBLE_ARRAY_MEMBER]; } WalSndCtlData; 其中SyncRepQueue是等待不同同步方式(备机日志写入磁盘、备机日志接收、备机日志回放等同步方式)的业务线程等待队列,用于当某一种同步方式满足条件之后,唤醒该类型的业务线程完成事务提交。lsn是上述几种队列队头后台线程等待的日志同步位置。sync_standbys_defined表示是否配置了同步备机。most_available_sync表示是否配置了最大可用模式;如果已配置,则在没有同步备机连接的情况下,后台业务线程可以直接提交,不用阻塞等待。sync_master_standalone表示当前是否有同步备机连接。demotion表示主机的降备方式。mutex表示保护walsnds结构体并发访问的互斥锁。walsnds表示保存wal sender的具体同步状态和进度信息。 3. 计划外切换(failover) 图4-40 failover流程示意图 如图4-40所示,failover(故障切换)时主机是异常状态,所以只有备机参与failover。failover的核心是让备机在满足一定条件以后退出日志复制和日志恢复流程。当数据库主线程“postmaster”线程(简称PM线程)在reaper中收到“startup”线程(即恢复线程)的停止信号后,将实例状态设置为PM_RUN,并将实例HA状态设置为PRIMARY_MODE。 4. 计划内切换(switchover) 图4-41 switchover流程示意图 如图4-41所示,switchover的过程比failover多了主机降备的处理,备机的流程和failover流程一致,因此没有在图中标出,参考failover流程即可。 5. 备机重建 图4-42 备机重建流程示意图 如图4-42所示,备机重建的过程相当于对主机进行了一次全量备份和恢复的操作,主要步骤包括:清理残留数据、全量拷贝数据文件、复制增量日志、启动备实例。这个过程中比较关键的两点是:文件和日志的拷贝顺序,以及备机第一次启动时选择的日志恢复起始位置。 6. cstore数据复制 在openGauss中,对于cstore表的数据复制与上述介绍略有不同。在一主多备部署场景下,每个CU填充写盘之后都会将CU整体数据记录到日志文件中,从而通过主备的日志复制和备机的日志回放,就可以实现cstore表增量数据的主备同步。在主备从部署场景下,每个CU填充写盘之后会直接将该CU数据拷贝到主机与备机之间的数据发送线程的局部内存中,并在事务提交之前阻塞等待数据发送线程传输完增量的CU数据才能完成事务提交,因此也实现了cstore表增量数据的主备同步。 下一篇我们详细讲述“4.3 内存表”相关内容。

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

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

上一篇我们详细讲述“3. astore元组多版本机制”、“4.astore访存管理”及“5.astore空间管理和回收”相关内容。本篇我们将继续为小伙伴们带来“4.2.4 ustore”的详细介绍。 4.2.4 ustore ustore属于In-place Update更新模式,中文意思为:原地更新,是openGauss内核新增的一种存储模式。openGauss内核当前使用的行引擎采用的是Append Update(追加更新)模式,该模式在INSERT、DELETE、HOT UPDATE(页面内更新)的场景下有较好的表现。但对于非HOT UPDATE场景,垃圾回收不够高效。 In-place Update存储模式提供“原地更新”能力,主要思路是将最新版本的“有效数据”和历史版本的“垃圾数据”分离存储。将最新版本的“有效数据”存储在数据页面上,而单独开辟一段undo(回滚)空间,用于统一管理历史版本的“垃圾数据”,因此数据空间不会由于频繁更新而膨胀,垃圾回收效率更高。通过NUMA-aware的undo子系统设计,使得undo子系统在多核平台上高效扩展。同时通过对元组和数据页面结构的重新设计,减少存储空间的占用。采用多版本索引技术,解决索引膨胀问题,彻底去除autovacuum(垃圾清理线程)机制,提升存储空间的回收复用效率。 1. 整体框架及代码概览 数据库中数据处理的本质是在保证ACID的基础上支持尽量高的并发查询。这种状况下,并发控制、页面多版本控制以及页面存储结构相互耦合在一起,数据库存储引擎需要进行整体设计从而在高并发的状况下保证各个事务处理看到类似串行执行的效果。 在整个技术体系中多版本控制用来提升读写并发能力,按照多版本排列方式可以分为两类。 (1)Oldest to New,即版本按照从最老到最新的方式进行链接,当一个事务访问该元组时,先看到这个元组最老的版本,同时使用对应的可见性判断机制,看是否是自己可见的版本,如果不是则沿着版本链条继续往后看较新的版本是不是自己需要的。 (2)Newest to Old,即版本按照从最新到最老的方式进行连接,当一个事务访问该元组时,先看到这个元组最新的版本,同时使用对应的可见性判断机制,看是否是自己可见的版本,如果不是则沿着版本链条继续往后看较老的版本是不是自己需要的。 在上面的描述中又引出一个设计点,如何组织新老数据,有如下几种方式。 (1) 将新数据和老数据放在同样的页面内,即每个数据页内放置着各个元组的新老数据,在需要进行不可见数据版本回收的时候需要遍历所有的页面。 (2) 将最新数据和老数据分离存储,在实际的数据页面内放置最新版本数据,所有的老版本数据都集中存储,新版本数据通过一个指针指向老版本所在的数据区域,当进行不可见老版本数据回收的时候只要扫描老版本集中存放的位置即可。 当新老数据分别存储的时候又引出第三个设计点,在对同一个页面或者元组反复读取时,是否要还原对应的页面在数据缓冲区中,这个设计点有如下几种方式。 (1) 访问旧元组所在的页面时,还原该页面,并将该页面的旧版本放入数据缓冲区中,节省一定时间内其他线程多次访问该版本页面带来的合成开销。弊端是占用更大的内存空间,同时缓冲区淘汰管理在原始LRU(Least Recently Used,最近最少使用算法)基础上同时要考虑页面版本。这种方式对应PCR(Page Consistency Read,页面一致性读),其本质的设计理念是空间换时间。 (2) 访问元组时,沿着版本链还原该元组,直到找到自己对应的版本。这种方式对于短时间访问冲突不高的场景能够降低内存使用,但如果短时间内高频访问一个页面内的元组,则每次都会遍历版本链造成访问效率低下。这种方式对应RCR(Row Consistency Read,行一致性读)。 按照上面的描述,整个多版本控制设计分为三个维度,如图4-10、表4-15所示。 图4-10 多版本控制设计维度 表4-15 多版本控制设计维度 维度 备选 版本存储方式 集中存储、分离存储 版本链组织方式 Oldest to New、Newest to Old 老版本管理方式 1、RCR,2、PCR 当前openGauss在版本存储方式、版本链组织方式上的设计选择是集中存储 + Oldest to New,在清理数据旧版本时需要遍历所有的页面找到不可见的元组版本然后清除。商用及开源的常见数据库的多版本控制设计三维度选择如表4-16所示。 表4-16 当前数据库多版本控制设计选择 数据库 架构设计选择 版本存储方式 版本链组织方式 老版本管理方式 常见数据库 分离存储 Newest to Old PCR 集中存储 Oldest to New PCR 分离存储 Oldest to New PCR 不同的多版本控制设计都不能做到尽善尽美,都有些不足之处,相关的缺点如下。 (1)多核系统上扩展性较差,不支持多核处理器的NUMA感知; (2)依赖于Vacuum进行老版本回收,后台线程定期清理; (3)缺乏对索引多版本,全局索引、闪回等功能的支持; (4)PCR管理方式,内存管理开销较大。 openGauss的ustore存储模式最大程度结合各种设计的优势,在多版本管理上的架构设计采取的组合如表4-17所示。 表4-17 ustore在多版本管理上的架构设计 维度 架构设计选择 版本存储方式 分离存储 版本链组织方式 Newest to old 老版本管理方式 PbRCR(Page Based RCR,基于页面的行一致性读) 同时为了事务能够跨存储格式查询,并复用现有备份、恢复、升级等能力,openGauss定义如下的融合引擎架构设计原则。 (1) 一套并发控制系统。 (2) 一套系统表管理系统。 (3) 一套日志管理系统。 (4) 一套锁管理系统。 (5) 一套恢复系统。 ustore架构如图4-11所示。 ustore和astore共用事务管理、并发控制、缓冲区管理、检查点、故障恢复管理与介质管理器管理。ustore主要功能模块如表4-18所示。 表4-18 ustore主要功能模块 模块 说明 代码位置 ustore表存取管理 向上对接SQL引擎,提供对ustore表的行级查询、插入、删除、修改等操作接口,向下根据ustore表页间、页内结构,以及ustore表元组结构,完成对ustore表文件的遍历和增删改查操作 主要在“src/gausskernel/storage/access/ustore”目录(单表文件管理)下 ustore索引存取管理 向上对接SQL引擎,提供对索引表的行级查询、插入、删除等操作接口,向下根据索引表页间、页内结构,以及索引表元组结构,完成对指定索引键的查找和增删操作 抽象框架代码在“src/gausskernel/storage/access/ubtree”目录下 ustore表页面结构 包括ustore表元组在页面内的具体组织形式,在页面内插入元组操作、页面整理操作、页面初始化操作等 主要代码在“src/gausskernel/storage/access/ustore/knl_upage”目录中 ustore表元组结构 包括ustore表元组的结构、填充、解构、修改、字段查询、变形等操作 主要代码在“src/gausskernel/storage/access/ustore/knl_utuple.cpp”文件中 Undo记录结构 包括undo记录的结构、填充、编码、解码等操作 主要代码在“src/gausskernel/storage/access/ustore/undo”目录中 多版本索引 包括 ustore 专用多版本索引 ubtree 的页面结构、查询、修改、可见性检查、垃圾回收等模块 主要代码在“src/gausskernel/storage/access/ubtree”目录中 2. 页面元组结构 1) 元组结构 本节介绍行存储引擎ustore表的页面元组结构。 元组结构的定义如下: typedef struct UHeapDiskTupleData { ShortTransactionId xid; uint16 td_id : 8, locker_td_id : 8; uint16 flag; uint16 flag2; uint8 t_hoff; uint8 data[FLEXIBLE_ARRAY_MEMBER]; } UHeapDiskTupleData; 该结构体只是元组头部的定义,真正的元组内容跟在该结构体之后,距离元组头部起始处的偏移由t_hoff成员保存。上面元组头部结构体部分成员信息同时也构成了该元组的系统字段(字段序号小于0的那些字段)。对各个结构体成员的含义说明如下。 (1) flag,元组属性掩码。包含是否有空字段标记、是否有外部TOAST标记、是否有变长字段标记、指定的事务槽位是否已被重复使用标记,以及更新、删除、锁等标记。 (2) flag2,元组另一个属性掩码。包含元组中字段个数。 (3) t_hoff,元组数据距离元组头部结构体起始位置的偏移。 (4) data,字段的NULL值bitmap,每个字段对应一个bit位,因此是变长数组。 ustore元组头部比astore元组头部小一半,因此在相同大小的页面上,ustore可以放置更多的元组。 在内存中,上述元组结构体使用时被嵌入在一个更大的元组数据结构体中,除了保存元组内容的disk_tuple成员之外,其他的成员保存了该元组的一些其他系统信息,并构成了该元组剩余的一些系统字段内容,定义如下: typedef struct UHeapTupleData { uint32 disk_tuple_size; uint1 tupTableType = UHEAP_TUPLE; uint1 tupInfo; int2 t_bucketId; ItemPointerData ctid; Oid table_oid; TransactionId t_xid_base; TransactionId t_multi_base; UHeapDiskTupleData* disk_tuple; } UHeapTupleData; 该结构体几个主要成员的含义如下。 (1) disk_tuple_size,元组长度。 (2) ctid,元组所在页面号和页面内元组指针下标。 (3) table_oid,该元组属主表的OID。 常用的元组操作接口和说明如表4-19所示。 表4-19 常用的元组操作接口 函数名 操作含义 UHeapFormTuple 利用传入的、各个元组字段的值数组,生成一条完整的元组,一般用于插入操作 UHeapDeformTuple 利用传入的完整元组以及各个字段的类型定义,解构各个字段的值,生成值数组,一般用于更新前的准备工作 UHeapFreetuple 释放一条元组对应的内存空间 UHeapCopyTuple 复制一条完整的元组,包括元组头和元组内容 UHeapSlotGetAttr 获取一条元组中指定的用户或系统字段值 UHeapGetSysAttr 获取一条元组中指定的系统字段值 UHeapCopyHeapTuple 从ustore槽位构造一条astore元组 UHeapToHeap 将一条ustore元组转换为一条astore元组 HeapToUHeap 将一条astore元组转换为一条ustore元组 2) 页面结构 ustore与astore相同,在openGauss中也使用默认的8kB页面。其结构如图4-12所示。 图4-12 ustore引擎页面结构示意图 在一个页面中,页面头部分对应的UHeapPageHeaderData结构体存储了整个页面的重要元信息。UHeapPageHeaderData之后有一个共享的页内事务目录(Transaction Directory,TD),对应元组指针变长数组。元组指针变长数组的每个数组成员存储了页面中从后往前的、每个元组的起始偏移和元组长度。如图4-12所示,真正的元组内容从页面尾部开始插入,向页面头部扩展;相应地,TD插槽目录与记录每条元组的元组指针从页面头定长成员之后插入,往页面尾部扩展。这样整个页面中间就会形成一个空洞,以供后续插入的元组和元组指针使用。每一个ustore表里的一条具体元组都有一个全局唯一的逻辑地址(和astore表里的元组相同),它由元组所在的页面号和页面内元组指针数组下标组成。 页面头具体结构体定义如下: typedef struct UHeapPageHeaderData { PageXLogRecPtr pd_lsn; uint16 pd_checksum; uint16 pd_flags; uint16 pd_lower; uint16 pd_upper; uint16 pd_special; uint16 pd_pagesize_version; uint16 potential_freespace; uint16 td_count; TransactionId pd_prune_xid; TransactionId pd_xid_base; TransactionId pd_multi_base; uint32 reserved; } UHeapPageHeaderData; 其中各个成员的含义如下。 (1) pd_lsn:该页面最后一次修改操作对应的预写日志位置的下一位,用于检查点推进和保持恢复操作的幂等性。 (2) pd_checksum:页面的CRC校验值。 (3) pd_flags:页面标记位,用于保存各类页面相关的辅助信息,如页面是否有空闲的元组指针、页面是否已满等。 (4) pd_lower:页面中间空洞的起始位置,即当前已使用的元组指针数组的尾部。 (5) pd_upper:页面中间空洞的结束位置,即下一个可以插入元组的起始位置。 (6) pd_special:页面尾部特殊区域的起始位置。该特殊位置位于第一条元组记录和页面结尾之间,用于存储一些变长的页面级元信息,如索引的辅助信息等。 (7) pd_pagesize_version:页面的大小和版本号。 (8) potential_freespace:页面中已被删除和更新的元组的潜在空间。 (9) td_count:共享的页内事务信息描述插槽的数量。 (10) pd_prune_xid:页面清理辅助事务号(64位),通常为该页面内现存最老的删除或更新操作的事务号,用于判断是否要触发页面级空闲空间整理。 (11) pd_xid_base:该页面内所有元组的基准事务号(64位)。该页面所有元组实际生效的64位XID事务号由pd_xid_base(64位)和元组头部的XID成员(32位)相加得到。 (12) pd_multi_base:类似pd_xid_base。当对元组加锁时,会将持锁的事务号写入元组中,该64位事务号由pd_multi_base(64位)和元组头部的XID(32位)相加得到。 页面的主要管理接口如表4-20所示。 表4-20 页面管理接口函数 函数名 操作含义 UPageInit 初始化一个新的ustore页面 UPageAddItem 在页面中插入一条新的元组 UHeapPagePruneOptPage 页面空闲空间整理 为了节省每个元组存储空间,元组头部UHeapDiskTupleData采用32位元组XID的组合设计方式。64位的pd_xid_base和pd_multi_base储存在页面上,元组上储存32位的XID。页面上pd_xid_base和pd_multi_base也需要通过额外的逻辑进行维护:同一个页面中所有元组实际的64位XID,一定要在pd_xid_base和pd_xid_base+232之间,所以如果新写入的事务号和页面上现有任意一个元组的XID事务号差距已经超过232,那么需要尝试对现有元组进行基线移位操作,更新pd_xid_base和pd_multi_base。 3)事务目录 事务目录是一种常用的共享资源。它可以为数据页上的元组(tuple)链接相应的事务表(Transaction Table)及undo子系统中的undo页面。数据库中的每个表可以自定义事务目录的数量,并可以复用那些已完成事务占据的事务目录。 每个数据页默认会有4个事务目录。根据并发需求的不同,事务目录的数量可设置为2到128之间的任意值。在使用CREATE TABLE命令创建表时添加了一个新的选项INIT_TD以声明所需的事务目录数量: CREATE TABLE t1 ( c1 integer; c2 boolean; ) WITH (INIT_TD=16); 当需要为新事务目录留位置时,系统会先查找当前页面中是否有空事务目录。若无空事务目录,系统将遍历事务目录列表来寻找可以复用的条目。条目是否可以复用取决于与该条目关联的事务的状态。 通常可以复用那些与已冻结或已中止的事务关联的事务目录。 (1) 对于已经冻结的XID,并复用该事务目录。 对于astore而言,冻结的XID代表着事务在所有的会话中都已经不再活跃。 而在ustore中,仅当一个事务创建的所有的回滚记录都被丢弃后,或者说没有其他的Snapshot需要再观察该事务创建的元组历史版本(tuple version)时,才将该XID视为冻结。ustore中的undo回收进程会维护一个oldestXidInUndo变量,系统将通过比较XID与该变量来确定XID是否含有回滚记录。如果XID < oldestXidInUndo,代表所有该XID产生的回滚记录都已经被丢弃。 (2) 对于已中止的事务,在该事务被回滚后,系统才会复用相应的事务目录条目。 (3) 对于已提交的事务,系统将不会无效化回滚记录地址,这样可以保证undo链的完整性。 当没有事务目录可以复用时,事务目录将会自动扩容以容纳更多的条目。需注意的是,事务目录的后面跟随着元组指针区,在扩展时,首先需要将row pointer array向右挪动来腾出空间。扩展后,新的事务目录条目将会在先前的事务目录条目之后依序添加。设计上,允许事务目录的容量最多扩至页面大小的约25%,即约100个事务目录(在8kB大小的页面中,约20Bytes/事务目录)。目前,系统将以每次增加两个事务目录的方式逐步扩容,最多扩至128个事务目录。ustore暂不支持收缩事务目录空间。 在扩容时,可以增加的总条目数也取决于当前页面中的可用空间。有时,页面中的总剩余空间并不能支持事务目录的扩容。此时若当前操作为INSERT或MULTI-INSERT,事务将会索取一个新的页面来进行操作。若操作为UPDATE或DELETE,事务将等待10毫秒后重试获取事务目录。Lock timeout设置可以控制获取事务目录的最大等待时间。在多由短事务组成的工作负载中,等待是可以接受的。 PG stats会报告事务目录等待等信息,以方便监测系统及描述工作负载。 事务目录申请的过程(UHeapPageReserveTransactionSlot函数)如图4-13所示。 图4-13 事务目录申请处理流程 如果当前事务需要申请一个新的事务目录,且系统中不存在空的事务目录时,系统会遍历所有事务目录并寻找可复用的事务目录。 (1) 首先系统会遍历事务目录,寻找XID < oldestXidInUndo的事务目录。这些条目将被视为已冻结。 (2) 接着系统会遍历目标页面上的元组。 ① 系统把已删除的元组标记为死亡,其余的标记为闲置。 ② 如果系统发现元组还在活跃状态,且相应的TD条目存在于步骤(1)给出的冻结列表之中,系统会把该事务目录设置为UHEAPTUP_SLOT_FROZEN(冻结)。 ③ 设置为冻结之后,事务目录中的XID及Undo指针会被无效化。 (3) 如果上述的冻结操作并未产生可用的槽位,系统会遍历事务目录并寻找与已提交或已中止事务关联的条目。这些条目在满足一定条件的状况下可被复用。 (4) 遍历目标页面上的元组。 ① 如果系统发现元组关联的事务目录存在于步骤(3)给出的已提交列表中,系统就把该TD条目的flag设为UHEAP_INVALID_XACT_SLOT(无效)。 ② 此外,这些事务目录的XID被重设为无效XID。但为了维护undo链的完整,undo指针将被保留。 (5) 如果并未找到与已提交事务关联的事务目录,最后将寻找与已中止事务关联的事务目录。 (6) 遍历与已中止事务关联的事务目录:对于每个事务目录,沿着undo链执行相关的undo操作。 (7) 如果并未找到事务目录,扩展事务目录。 (8) 返回结果。 3. 回滚段设计与MVCC 1) 回滚段 旧版本数据会集中在回滚段的undo目录中,为了减少读写冲突,旧版本数据(回滚记录)采用追加写的方式写入数据目录的undo目录下。这样旧版本数据的读取和写入不会发生冲突,同一个事务的旧版本数据也会连续存放,便于进行回滚操作。为了减少并发写入时的竞争,undo目录空间被划分成多个逻辑区域(UndoZone,回滚段逻辑区域)。线程会在自己的逻辑区域上进行分配,与其他线程完全隔离,从而写入旧数据分配空间时就不会有额外的锁开销。UndoZone还可以按照CPU的NUMA核进行划分,每个线程会从当前的NUMA核上的UndoZone进行分配,进一步提升分配效率。在分配undo空间时会按照事务粒度进行记录,旧版本数据一旦确认没有事务进行访问,就会进行回收。 为了在回滚段的空间寻址,回滚记录使用8字节的指针来进行寻址,如图4-14所示。 图4-14 回滚记录寻址指针 其中各个字段的含义如下: (1) zoneId:占用20bit,表示逻辑区域的ID。 (2) blockId:占用31bit,表示块号,默认为8k。 (3) offset:占用13bit,表示块内偏移。 旧版本的数据采用回滚记录的格式存入回滚段中,其中回滚记录的格式如下所示: Class UndoRecord { … UndoRecordHeader whdr_; UndoRecordBlock wblk_; UndoRecordTransaction wtxn_; UndoRecordPayload wpay_; UndoRecordOldTd wtd_; UndoRecordPartition wpart_; UndoRecordTablespace wtspc_; StringInfoData rawdata_; } 其中,除了rawdata_代表了旧版本数据,其他成员均为结构体,下面对每个结构体分别进行说明。 whdr_成员由下面的结构组成: typedef struct { TransactionId xid; CommandId cid; Oid reloid; Oid relfilenode; uint8 utype; uint8 uinfo; } UndoRecordHeader; 各个字段的含义如下。 (1) xid:生成此回滚记录的事务ID,用于检查事务的可见性。“2)MVCC”小节有介绍。 (2) CID(Command ID,命令ID):生成此回滚记录的命令ID,用于判断可见性。 (3) reloid:relation对象的ID,回滚时需要。 (4) relfilenode:relfilenode对象的ID,回滚时需要。 (5) utype:操作类型,像UNDO_INSERT、UNDO_DELETE、UNDO_UPDATE等。 (6) uinfo:控制字段,用来判断后续的结构是否存在,用来减少回滚记录的占用空间。 wblk_成员由下面的结构组成: typedef struct { UndoRecPtr blkprev; BlockNumber blkno; OffsetNumber offset; } UndoRecordBlock; (1) blkprev:指向同一个block前一条回滚记录,用于回滚和事务可见性。“2)MVCC”小节有介绍。 (2) blkno:block number(块号)。 (3) Offset:修改的tuple在row pointer中的偏移。 wtxn_成员由下面的结构组成。 typedef struct { UndoRecPtr prevurp; } UndoRecordTransaction; prevurp:当一个事务的回滚记录跨越两个UndoZone时,后续的回滚记录使用此指针指向前一条回滚记录。 wpay_成员由下面的结构组成。 typedef struct { UndoRecordSize payloadlen; } payloadlen:rawdata_的长度。 wtd_成员由下面的结构组成。 typedef struct { TransactionId oldxactid; } UndoRecordOldTd; oldxactid:旧版本数据里事务目录的事务ID。 wpart_成员由下面的结构组成。 typedef struct { Oid partitionoid; } UndoRecordPartition; partitionoid:分区表的分区对象OID。 wtspc_成员由下面的结构组成。 typedef struct { Oid tablespace; } UndoRecordTablespace; tablespace:表空间的OID。 回滚段使用事务目录来记录每个事务分配的undo空间,便于事务回滚和回收。事务发生回滚时,会读取事务目录中记录的undo空间的起始位置,再读取undo空间中的回滚记录进行回滚操作,其中回滚记录中的字段如下: class TransactionSlot { TransactionId xactId_; UndoRecPtr startUndoPtr_;/*事务分配的undo空间开始*/ UndoRecPtr endUndoPtr_;/*事务分配的undo空间结束*/ uint8 info_;/*标记:如事务回滚状态*/ Oid dbId_;/*数据库对象ID*/ } (1) xactId:事务ID。 (2) startUndoPtr:事务分配的undo空间开始位置。 (3) endUndoPtr:事务分配的undo空间结束位置。 (4) info_:标记值,如事务回滚状态。 (5) dbId:数据库对象ID。 回滚段提供分配undo空间和更新事务目录的接口,主要接口如表4-21所示。 表4-21 回滚段主要接口 接口名 含义 AllocateUndoSpace 为回滚记录分配undo空间 UpdateTransactionSlot 更新事务目录 以ustore的删除操作为例,undo空间分配流程如下。 (1) UheapDelete作为ustore的删除接口,会调用UHeapPrepareUndoDelete函数准备回滚记录(undo record)。UHeapPrepareUndoDelete函数会填充回滚记录的各个字段(其中旧数据会设置到回滚记录的raw data字段上),再调用PrepareUndoRecord函数分配undo空间。PrepareUndoRecord函数调用“undo::AllocateUndoSpace”函数分配undo空间,再读取对应的回滚记录到缓冲池中。AllocateUndoSpace函数不仅会为回滚记录分配空间(使用“UndoZone::AllocateSpace”函数),如果是事务的第一条回滚记录,还会调用“UndoZone::AllocateSlotSpace”函数为事务目录分配空间。AllocateSpace函数会进行判断,如果回滚记录超过当前undo file的大小,就扩展当前的undo file,AllocateSlotSpace函数的逻辑类似。 (2) UheapDelete函数调用InsertPreparedUndo函数,将准备好的回滚记录追加写到缓冲池中的回滚段页面。 (3) UheapDelete函数调用UpdateTransactionSlot,记录下该事务分配的undo空间起始、事务ID、数据库ID。如果是事务的第一次更新,会从事务目录空间分配新的事务目录再进行更新。 undo空间需要回收回滚记录来保证undo空间不会无限膨胀,一旦事务id小于当前快照中最小的Xmin(oldestXmin),回滚记录中的旧版本数据就不会被访问,此时就可以对回滚记录进行回收。 如前述描述undo空间中的回滚记录按照事务ID递增的顺序存放在UndoZone中,回收的条件如下所示。 (1) 事务已经提交并且小于oldestXmin的undo空间可以回收。 (2) 事务发生回滚但已经完成回滚的undo空间可以回收。 图4-15 undo回收过程 如图4-15所示,UndoZone1中回收到小于oldestXmin的已提交事务16068,UndoZone2中回收到16050,UndoZone m回收到16056。而UndoZone n回收到事务16012,而事务16014待回滚但未发生回滚,因此UndoZone n回收事务id上限只到16014。其他zone的上限是oldestXmin,oldestXidInUndo会取所有undozone上的上限最小值,因此oldestXidInUndo等于16014。undo回收主要函数如表4-22所示。 表4-22 undo回收主要函数 函数名 操作含义 UndoRecycleMain 回收线程的入口函数,会在每个zone上调用RecycleUndoSpace函数 RecycleUndoSpace 按照前述条件回收undo空间,记录日志 2) MVCC ustore的可见性检查和astore类似,将快照CSN和元组删除和插入事务的CSN进行比较,判断元组是否可见。ustore和astore使用同一套事务管理机制和快照管理机制。 ustore和astore最大的区别在于astore会在页面上保留旧版本数据,而ustore在将旧版本数据放到回滚段统一存放。在需要获取旧版本数据时,astore可以直接从tuple的头部读取到元组的插入和删除的事务号(XID),来判断元组的可见性。但是ustore需要从回滚段里读取旧版本的事务信息,来判断旧版本是否可见。由于从回滚段中读取旧版本数据存在相对昂贵的开销,ustore通过一系列的优化手段来避免从回滚段中读取旧版本数据。 ustore在获取元组时,会先检查对应的事务目录。事务目录分成有效和无效两种。当事务目录是有效的,ustore直接就会得到元组上最新的事务。 如果事务目录被冻结(FROZEN),意味着元组已经在所有的事务中都会可见。如果事务目录中的事务id小于oldestXidInUndo,意味着元组已经足够旧在所有事务中都可见。同时会把事务目录置成冻结,来加速后续的查询。 如果元组被标记有一个无效事务目录,意味着修改元组的事务已经提交,并且比当前的事务目录中的事务旧。此时ustore会使用事务目录中的事务进行可见性判断。如果可见,意味着修改元组的事务更已经可见,就不需要从undo目录中再读取事务信息。 图4-16 元组查询过程 元组不可见的场景,ustore会从undo目录中读取回滚记录中的旧版本数据查找元组。例子如图4-16所示。查找tbl表中c1=1的数据项,从索引中读取到数据项位于block 1和offset 2,使用UHeapTupleFetch函数再从block 1中查询到元组,需要判断该元组的可见性。 (1) 从元组的TD读到ITL2,和astore类似,根据CSN的大小,判断TD2中的XID不可见,需要使用GetTupleFromUndo函数读取回滚记录。 (2) GetTupleFromUndo函数调用GetTupleFromUndoRecord函数读取回滚记录,使用InplaceSatisfyUndoRecord函数判断其中的block 1和offset 2是满足要求的元组。但是XID=1610可以判断出当前页面的tuple不可见,ustore继续查询更老的版本。由于旧元组的TD 1和当前的TD 2不一致,使用UHeapUpdateTDInfo从TD 2 undo链条进行切换,根据旧元组的TD 1找到当前的undo指针找到前一次修改。 (3) 再次读取到回滚记录,其中的block 1和offset 1并非要找的元组,ustore继续查询更老的版本,根据blkprev指针读取前一次修改。 (4) 读取到回滚记录,其中的block 1和offset 3并非要找的元组,ustore继续查询更老的版本,根据blkprev指针读取前一次修改。 (5) 读取到回滚记录,其中的block 1和offset 2是要求的元组,ustore判断可见性。根据CSN的大小,事务可见。因此前一次命中的元组可见,即(1, abc)可见,因此查找到元组的c2等于abc。 4. 多版本索引 在openGauss中实现了多版本索引ubtree,是专用于ustore的B-Tree索引变种,相比原有的B-Tree索引有如下差异点。 (1) 支持索引数据的多版本管理及可见性检查,能够自主鉴别旧版本元组并进行回收,同时索引层的可见性检查使得索引扫描(Index Scan)及仅索引扫描(Index Only Scan)性能有所提升。 (2) 在索引插入操作之外,增加了索引删除操作,用于对被删除或修改的元组对应的索引元组进行标记。 (3) 索引按照key + TID的顺序排列,索引列相同的元组按照对应元组的TID作为第二关键字进行排序。 (4) 添加新的可选页面分裂策略“insertpt”。 ubtree实现了索引访问接口所要求的全部接口,如表4-23所示: 表4-23 ubtree访问接口函数 接口名称 对应函数 接口含义 aminsert ubtinsert 插入一个索引元组 ambeginscan ubtbeginscan 开始一次索引扫描 amgettuple ubtgettuple 获取一个索引元组 amgetbitmap ubtgetbitmap 通过索引扫描获取所有元组 amrescan ubtrescan 重新开始一次索引扫描 amendscan ubtendscan 结束索引扫描 ammarkpos ubtmarkpos 标记一个扫描位置 amrestpos ubtrestpos 恢复到一个扫描位置 ammerge ubtmerge 合并多个索引 ambuild ubtbuild 建立一个索引 ambuildempty ubtbuildempty 建立一个空索引 ambulkdelete ubtbulkdelete 批量删除索引元组 amvacuumcleanup ubtvacuumcleanup 索引后置清理 amcanreturn ubtcanreturn 是否支持 Index Only Scan amcostestimate ubtcostestimate 索引扫描代价估计 amoptions ubtoptions 索引选项 此外,还实现了新增的的索引删除函数UBTreeDelete。 1) 索引页面组织 多版本索引层次结构与B-Tree索引基本相同,非叶子节点与B-Tree索引保持一致,仅页尾的Special字段有所不同。ubtree中的Special字段UBTPageOpaqueDataInternal如下所示: typedef struct UBTPageOpaqueDataInternal { …… /* 以上部分与BTPageOpaqueDataInternal一致 */ TransactionId last_delete_xid; /* 记录页面上最后一次删除事务的 XID */ TransactionId xid_base; /* 页面上的 xid-base */ int16 activeTupleCount; /* 页面上活跃元组计数 */ } UBTPageOpaqueDataInternal; typedef UBTPageOpaqueDataInternal* UBTPageOpaqueInternal; 其中last_delete_xid与activeTupleCount用于索引的自治式回收,会在ustore中的“6. 空间管理和回收”一节详细讲解。 通过xid_base字段,页面上的XID可以仅储存基于该xid_base的一个32位偏移(Offset),节省XID存储的空间开销。实际的XID为页面上的xid_base加上存储的XID(也就是偏移)得到。 多版本中的叶子页面的结构如图4-17所示。 图4-17 ubtree 叶子页面结构示意图 与astore堆页面中维护版本信息的方法类似,ubtree的叶子节点中每个索引元组尾部都附加了对应的xmin和xmax。由于索引只是用于加速搜索的结构,本身不与历史版本概念强相关,仅通过xmin和xmax来标识这个索引元组是从什么时候开始有效的,又是什么时候被删除的,而不像astore中堆元组一样会有指向旧版本元组的指针。 新插入的索引元组尾部用于存放xmin和xmax 空间在ubtinsert函数执行的过程中预留出来。预留的空间及xmin在索引元组插入时通过UBTreePageAddTuple函数中写入页面,而xmax在索引元组删除时通过UBTreeDeleteOnPage函数中写入页面。 在UBTreePagePruneOpt函数中,索引元组通过其xmin和xmax信息来判断该元组是否已经无效(Dead),进而进行独立的页面清理。该函数会尝试清除所有无效的元组,并进行相应的碎片整理。 索引扫描时会调用UBTreeFirst函数定位到第一个满足扫描条件的索引元组,然后调用UBTreeReadPage获取当前页面中符合索引扫描条件,且能够通过可见性检查的元组。可见性检查通过UBTreeVisibilityCheckXid函数及UBTreeVisibilityCheckCid函数处理,其基本逻辑与astore类似,通过xmin与xmax及当前的快照进行可见性判断。 在ubtree中,索引元组除了按照索引列有序排列之外,对于索引列相同的元组,还将其对应堆元组的TID作为第二关键字进行排序。其具体实现大致都集中在ubtbuild函数及ubtinsert函数调用的过程中,这中间对索引列相同的元组会按照TID来进行额外的比较。实现还借助了BTScanInsert结构体,该结构体定义如下: typedef struct BTScanInsertData { bool heapkeyspace; /* 标志索引是否额外按 TID 排序 */ bool anynullkeys; /* 标志待查找的索引元组是否有为 NULL的列 */ bool nextkey; /* 标志是否希望寻找第一个大于扫描条件的元组 */ bool pivotsearch; /* 标志是否希望查找 Pivot 元组 */ ItemPointer scantid; /* 用于作为排序依据的 TID */ int keysz; /* scankeys 数组的大小 */ ScanKeyData scankeys[INDEX_MAX_KEYS]; } BTScanInsertData; 在索引元组将TID作为第二关键字排序之后,用于划分搜索空间的非叶子节点元组及叶子节点的Hikey元组(统称Pivot元组)也需要携带对应的TID信息。这会使得Pivot元组占用空间增加,非叶子的扇出(fan out)降低。为了避免这一特性导致的扇出降低,若不需要比较TID即可区分两个叶子页面,则对应的Pivot原则中就不需要储存TID信息。类似地,Pivot元组中也可以去掉一些不需要进行比较的索引列,这一逻辑在UBTreeTruncate函数中进行处理。原则是当比较前几列就可以区分两个叶子页面时,Pivot元组中就不需要储存后续的列。 2) 索引操作 对于原有的B-Tree索引而言,主要有四类操作:索引创建、索引扫描、索引插入以及索引删除。下面将依次进行介绍。 (1) 索引创建。 索引创建操作由索引上的ubtbuild函数及ustore上的IndexBuildUHeapScan函数配合完成。IndexBuildUHeapScan函数负责扫描对应的ustore表,并取出每个元组的最新版本(遵循SnapshotNow的语义)以及其对应的xmin和xmax。若发现某个元组存在被就地更新的旧版本,则会将该索引标记为HotChainBroken。被标记为HotChainBroken的索引,会复用astore原有的逻辑,禁止隔离级别为可重复读(Read Repeatable)的老事务访问。ubtbuild函数会接收IndexBuildUHeapScan传过来的元组,将其按照索引列及TID排序后依次插入到索引页面中,并构建相应的元页面及上层页面。整个创建流程需要将所有页面都记录到XLOG中,并强制将存储管理中的内容刷到永久存储介质后才算成功结束。 (2) 索引扫描。 索引扫描与B-Tree索引基本一致,但是需要对索引元组进行可见性检查。没有通过可见性检查的索引元组不会被返回,通过可见性检查的元组仍需要在ustor 堆表上进行可见性检查,并找到正确的可见版本。在IndexOnlyScan场景中,通过可见性检查的元组即可直接返回,不需要再访问堆表。 不过索引进行可见性检查时,由于索引元组只存放了xmin和xmax而没有CID(对应“4.2.3 astore”节堆表元组中的t_cid字段)信息,如果发现了当前事务修改过的索引元组则不能正确地通过CID来判断其可见性。此时会将该元组视为可见,但会标记xs_recheck_itup,告知ustore的数据页面需要在取到对应的数据元组后,再次构建对应的索引元组并与返回的索引元组进行比较,来确认该索引元组是不是真正可见。相关逻辑在 UBTreeVisibilityCheckXid、UBTreeVisibilityCheckCid以及RecheckIndexTuple函数中进行处理。 (3) 索引插入。 索引元组需要存储对应的xmin和xmax版本信息,但其所占用的空间并不表现在IndexTupleSize中,而是对外部透明。索引插入的接口函数为ubtinsert,为了正确插入带有版本信息的元组,需要在执行插入前增加IndexTupleSize以预留用于储存版本信息的空间。真正将元组插入到页面的时候,会将版本信息所占用的空间大小从IndexTupleSize中去除。 在索引插入的过程中若页面空间不足,会首先调用UBTreePagePruneOpt函数尝试对已经无效的元组进行清理。若清理失败或清理成功后空间仍然不足,会进行索引页面分裂。索引页面分裂会在UBTreeInsertOnPage函数中进行。ubtree中存在两种分裂策略:default以及insertpt。其中default策略会将原页面上的内容均匀地分配到两个页面上,而insertpt会根据新插入元组的插入规律、插入位置及TID等信息选择合适的分裂点。 在ubtree需要申请新的页面时,并不会像原有的B-Tree索引一样调用_bt_getbuf通过FSM来查找可用页面。ubtree带有自治式的空间管理机制,通过UBtreeGetNewPage函数获取新页面。该自治式空间管理机制将在空间管理和回收部分介绍。 (4) 索引删除。 索引删除操作用于在堆元组被删除的同时,将对应的索引元组也标上对应的xmax。索引删除的流程与插入类似,通过二分查找定位到待删除元组的位置,并将xmax写入到对应的位置。需要注意的是,要删除的元组是索引列以及TID都匹配,且还未被写入xmax的那个元组,这部分逻辑在UBTreeFindDeleteLoc函数中处理。在最后会调用UBTreeDeleteOnPage函数为对应的索引元组写上xmax,更新页面上的last_delete_xid以及activeTupleCount,并在检测到activeTupleCount为0时将该页面放入潜在空页队列(Potential Empty Page Queue)中。关于潜在空页队列会在空间管理和回收部分介绍。 5. 存取管理 openGauss中的ustore表访存接口如表4-24所示。由于openGauss中ustore表只有一种页面和元组结构,因此在上述接口中,直接实现了底层的页面和元组操作流程。 表4-24 ustore表访存接口 函数名称 接口含义 heap_open 打开一个ustore表,得到表的相关元信息 heap_close 关闭一个ustore表,释放该表的加锁或引用 UHeapRescan 重新开始ustore表(顺序)扫描操作 UHeapGetNext (顺序)获取下一条元组 UHeapGetTupleFromPage UHeapGetNext内部实现,单页校验模式 UHeapScanGetTuple UHeapGetNext内部实现,单条校验模式 UHeapGetPage (顺序)获取并扫描下一个ustore表页面 UHeapInsert 在ustore表中插入一条元组 UHeapMultiInsert 在ustore表中批量插入多条元组 UHeapDelete 在ustore表中删除一条元组 UHeapUpdate 在ustore表中更新一条元组 UHeapLockTuple 在ustore表中对一条元组加锁 6. 空间管理和回收 不同于astore的空间管理和回收机制,ustore实现了自治式的空间管理机制。ustore里堆以及索引的空间分配和回收都在业务运行的过程中平稳地进行,不依赖中量级的VACUUM及AUTOVACUUM清理机制。 1) 自治式堆页面空间管理 ustore中堆页面的自治式空间管理,建立在与astore类似的轻量级堆页面清理机制的基础上。在执行DML及DQL操作的过程中,ustore都会进行堆数据页面清理,以取代VACUUM清理机制。UHeapPagePruneOptPage函数是页面清理的入口函数,会清理已经提交的被删除元组。 对于astore而言,复用数据元组的行指针前必须保证对应的索引元组已经被清理。这是为了防止通过索引元组访问已经被复用的行指针,导致取到错误的数据。在astore中需要通过VACUUM操作将这样的无效索引元组统一清除掉后才能复用行指针,这使得堆页面和索引页面的清理逻辑耦合在一起,也会导致间断性的大量I/O。在ustore中能高效地单独进行数据和索引页面的清理,因为带有版本信息的ubtree能够独立检测并过滤掉无效的索引元组,不会通过无效索引元组访问对应的数据表。 堆页面的空间管理机制复用openGauss中的FSM来管理UHeap中的可用空间。在UHeapPagePruneOptPage函数成功对页面进行清理后,会将其空闲空间刷新到对应的FSM页面中。为了避免每次页面清理都需要更新整个树状结构的FSM,从而带来额外的开销,引入了一个更新整个FSM的概率计算。考虑当前清理后的可用空间占预留可用空间(Reserved Free Space)阈值的百分比,计算得出清理一个页面后调用FreeSpaceMapVacuum函数的概率。也就是说,页面清理获得的可用空间越大,更新整个FSM的概率也就越大。 当数据元组被删除时,会在页面上记录对应的潜在空闲空间(Potential Free Space),该值用于估计页面上的空闲空间。在运行过程中,有多个场景会调用UHeapPagePruneOpt对页面尝试进行清理。DML语句执行过程中,INSERT、UPDATE以及DELETE操作都会拿到页面的写锁。如果发现空间不足,或者检测到潜在空闲空间到达某个阈值,会尝试对页面进行清理。DQL查询语句执行的过程中若检测到页面上潜在空闲空间到达阈值,也同样会尝试申请页面的写锁;如果拿到了页面的写锁,会尝试对页面进行清理。 存在可清理的元组,但一直不被访问的页面不能通过这一机制正确地清理。为了解决这一问题,引入了基于概率的清理方案。在RelationGetBufferForUTuple函数寻找新的可用空间时,若通过FSM发现没有足够的可用空间,在对物理文件进行扩展前,会“随机”选取一些页面进行清理。该机制并非完全随机选取,在多次尝试后选取的页面会覆盖到整个关系的全部页面。为了性能考虑,该过程中默认最多选取10个页面进行清理,该数量可以通过GUC参数max_search_length_for_prune进行设置。具体的页面选取数量通过DeadTupleRatio以及PruneSuccessRatio计算得出。其中DeadTupleRatio表示该表中无效元组的大致比例,该变量以统计信息的方式进行收集,在进行DML的过程中会进行更新;PruneSuccessRatio大致表示近几次尝试清理的成功率。 2) 自治式索引页面空间管理 索引页面的空间管理不依靠FSM数据结构,而是依靠特有的URQ(UBtree Recycle Queue)结构,简称为回收队列。索引回收队列单独储存在ubtree索引对应的.urq文件中,没有原有B-Tree索引的.fsm文件。索引回收队列相关代码在“ubtrecycle.cpp”文件中。涉及到的主要函数接口见表4-25。 表4-25 索引回收队列主要接口 函数名称 接口含义 UBTreeTryRecycleEmptyPage 尝试从潜在空页队列回收一个页面 UBTreeGetAvailablePage 获取有效页面(潜在空页或空闲页面) UBTreeRecordUsedPage 记录被成功使用的页面 UBTreeRecordEmptyPage 记录潜在的空页 UBTreeGetNewPage 获取新的可用页面 索引中的回收队列分为两部分,一部分是潜在空页队列(Potential Empty Page Queue),一部分是可用页面队列(Available Page Queue)。两个队列都是跨页面的循环队列,其中每个元素都会储存blkno以及XID。其中blkno表示该元素对应索引页面的block number;XID表示该页面在哪个时刻能够被回收或复用。这些元素在循环队列单个页内按照XID的顺序进行排序,以便于快速找到XID 小(最可能被回收或复用)的页面。其结构如图4-18所示。 图4-18 ubtree回收队列结构示意图 对于潜在空页队列而言,里面存放页内元组已经被全部删除但还没有全部无效的页面,其中的XID就标志页面中最后一个元组无效的可能时机。在系统整体的oldestXmin超过该XID后,该页面就有可能被从索引上删除,但也可能因为新插入元组或删除元组的事务中止而导致页面不能被删除。潜在空页队列中的页面在成功被删除后会被放入可用页面队列,并记录删除时最新事务的XID。 对于可用页面队列而言,里面存放已经被删除,可以或即将可以被复用的页面。其中XID就表示该页面可以被复用的时机。这样的页面复用时延是来自B-Tree索引页面删除时可能的并发访问导致的,可以参考nbtree文件夹下README 关于页面删除的部分。 在ubtree进行索引删除时,会更新页面上的last_delete_xid字段以及activeTupleCount字段。若更新后activeTupleCount变为0,会将该页面放入潜在空页队列,并将此时的last_delete_xid作为对应的可回收时间点。 在业务运行的过程中,索引会通过UBTreeTryRecycleEmptyPage函数不断尝试对潜在空页队列中的页面进行回收。在索引申请新的页面时,会通过UBTreeGetNewPage函数与可用页面队列交互,查找当前可用的空闲页面。当可用页面队列中没有可用页面时,一般会通过扩展索引物理文件的方式来获得新的页面。但也存在物理文件批量扩展,或扩展后还未来得及使用就出错退出的情况。此时在回收队列的元信息页面中保存了已正确追踪的页面数量,若该数量少于整个索引表的页面数量,会尝试去使用这一部分未追踪的页面,并更新已追踪的页面数量。 3) 中量级和重量级手动页面清理 与astore相同,ustore也提供VACUUM语句来让用户主动执行对某个ustore表及其上的索引进行中量级清理。其对外表现与astore一致,可参考astore的空间管理和回收内容。 在ustore中,中量级清理同样通过lazy_vacuum_rel函数进入,但不会调用lazy_scan_heap,而是调用LazyScanUHeap函数来进行数据页面的清理。在进行索引清理时,会调用lazy_vacuum_index接口及LazyVacuumHeap函数来清理索引文件和堆表文件,索引清理时会调用ubtbulkdelete函数。 重量级的VACUUM FULL也与astore一致,会清理无效数据并对数据空间和索引空间重新进行组织。重量级清理的对外接口是cluster_rel函数,本质上是重新对数据进行聚簇,清理过程中会阻塞对该表的所有操作。 介绍完“4.2.4 ustore”,下篇我们将详细介绍“4.2.5 行存储索引机制”相关内容,敬请期待!

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

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

上一篇我们将详细介绍“4.2.3 astore”相关内容,本篇我们将讲述“3. astore元组多版本机制”、“4.astore访存管理”及“5.astore空间管理和回收”。 4.2.3 astore 3. astore元组多版本机制 openGauss行存储表支持多版本元组机制,即为同一条记录保留多个历史版本的物理元组以解决对同一条记录的读、写并发冲突(读事务和写事务工作在不同版本的物理元组上)。 astore存储格式为追加写优化设计,其多版本元组产生和存储方式如图4-5所示。当一个更新操作将v0版本元组更新为v1版本元组之后,如果v0版本元组所在页面仍然有空闲空间,则直接在该页面内插入更新后的v1版本元组,并将v0版本的元组指针指向v1版本的元组指针。在这个过程中,新版本元组以追加写的方式和被更新的老版本元组混合存放,这样可以减少更新操作的I/O开销。然而,需要指出的是,由于新、老版本元组是混合存放的,因此在清理老版本元组时需要的清理开销会比较大。因此,astore存储格式比较适合频繁插入、少量更新的业务场景。 图4-5 astore多版本元组产生和存储方式示意图 下面结合图4-6,介绍openGauss中行存储格式多版本元组的运行机制: 图4-6 行存储格式多版本元组运行机制示意图 (1) 首先事务号为10的事务插入一条值为value1的新记录。对应的页面修改为:在0号物理页面的第一个元组指针指向位置,插入一条“xmin”字段为10、“xmax”字段为0、“ctid”字段为(0,1)、“data”字段为value1的物理元组。该事务提交,将CSN从3推进到4,并且在CSN日志中对应事务号10的槽位处记下该CSN的值。 (2) 然后事务号为12的事务将上面这条记录的值从value1修改为value2。对应的页面修改为:在0号物理页面的第二个元组指针指向位置,插入另一条“xmin”字段为12、“xmax”字段为0、“ctid”字段为(0,2)、“data”为value2的物理元组。同时保留上面第一条插入的物理元组,但是将其“xmax”字段从0修改为12,将其“ctid”字段修改为(0,2),即新版本元组的物理位置。该事务提交,将CSN从7推进到8,并且在CSN日志中对应事务号12的槽位处记下该CSN的值。 (3) 最后事务号为15的事务将上面这条记录的值从value2又修改为value3,对应的页面修改为:(假设0号页面已满)在1号物理页面的第一个元组指针指向位置,插入一条“xmin”字段为15、“xmax”字段为0、“ctid”字段为(1,1)、“data”字段为value3的物理元组;同时,保留上面第1、第2条插入的物理元组,但是将第2条物理元组的“xmax”字段从0修改为15,将其“ctid”字段修改为(1,1),即最新版本元组的物理位置。该事务提交,将CSN从9推进到10,并且在CSN日志中对应事务号15的槽位处记下该CSN的值。 (4) 对于并发的读事务,其在查询执行开始时,会获取当前的全局CSN值作为查询的快照CSN。对于上面同一条记录的3个版本的物理元组来说,该读查询操作只能看到同时满足如下两个条件的这个物理元组版本。 元组“xmin”字段对应的CSN值小于等于读查询的快照CSN。 元组“xmax”字段为0,或者元组“xmax”字段对应的CSN值大于读查询的快照CSN。 比如,若并发读查询的快照CSN为8,那么这条查询将看到value2这条物理元组;若并发读查询的快照CSN为11,那么这条查询将看到value3这条物理元组。 对于不同的行存储子格式,上述多版本元组的格式和存储方式可能有所不同,但是可见性判断和并发控制方式都是如图4-6中所示的。通过上面介绍的元组可见性判断流程,可以发现:并发的读事务会根据自己的查询快照在同一个记录的多个历史版本元组中选择合适的那个来返回。并且即使是在可重复读的事务隔离级别下,只要使用相同的快照总可以筛选出相同的那个历史版本元组。在整个过程中读事务不阻塞任何对该记录的并发写操作(更新和删除)。 更详细的元组可见性判断流程将在第5章中详细介绍。 最后,对于astore行存储格式,更新一条记录的详细执行流程如图4-7所示,该图可以帮助读者更形象地理解多版本元组的产生流程,以及写、写并发下的处理逻辑。 图4-7 更新astore记录的执行流程示意图 4. astore访存管理 openGauss中的astore堆表访存接口如表4-13所示。 表4-13 astore堆表访存接口 接口名称 接口含义 对应的行存储统一访存接口 heap_open 打开一个表,得到表的相关元信息 无 heap_close 关闭一个表,释放该表的加锁或引用 无 heap_beginscan 初始化堆表(顺序)扫描操作 tableam_scan_begin heap_endscan 结束并释放堆表(顺序)扫描操作 tableam_scan_end heap_rescan 重新开始堆表(顺序)扫描操作 tableam_scan_rescan heap_getnext (顺序)获取下一条元组 tableam_scan_getnexttuple heap_markpos 记录当前扫描位置 tableam_scan_markpos heap_restrpos 重置扫描位置 tableam_scan_restrpos heapgettup_pagemode heap_getnext内部实现,单页校验模式 无 heapgettup heap_getnext内部实现,单条校验模式 无 heapgetpage (顺序)获取并扫描下一个堆表页面 tableam_scan_getpage heap_init_parallel_seqscan 初始化并行堆表(顺序)扫描操作 tableam_scan_init_parallel_seqscan heap_insert 在堆表中插入一条元组 tableam_tuple_insert heap_multi_insert 在堆表中批量插入多条元组 tableam_tuple_multi_insert heap_delete 在堆表中删除一条元组 tableam_tuple_delete heap_update 在堆表中更新一条元组 tableam_tuple_update heap_lock_tuple 在堆表中对一条元组加锁 tableam_tuple_lock heap_inplace_update 在堆表中(就地)更新一条元组 无 以astore堆表顺序扫描为例,执行流程如下。 (1) 调用heap_open接口打开待扫描的堆表,获取表的相关元信息,如表的行存储子格式为astore格式等。该步通常要获取AccessShare一级表锁,防止并发的DDL操作。 (2) 调用tableam_scan_begin接口,从g_tableam_routines数组中找到astore的初始化扫描接口,即heap_beginscan接口,完成初始化顺序扫描操作相关的结构体。 (3) 循环调用tableam_scan_getnexttuple接口,从g_tableam_routines数组中找到astore的扫描元组接口,即heap_getnext接口,顺序获取一条astore元组,直到完成全部扫描。顺序扫描时,每次先获取下一个页面,然后依次返回该页面上的每一条元组。这里提供了两种元组可见性的判断时机: a) heapgettup_pagemode。在第一次加载下一个页面时,加上页面共享锁,完成对页面上所有元组的可见性判断,然后将可见的元组位置保存起来,释放页面共享锁。后面每次直接从保存的可见性元组列表中返回下一条可见的元组,无须再对页面加共享,使用快照的查询,默认都使用该批量模式,因为元组的可见性在同一个快照中不会再发生变化。 b) heapgetpage。除了第一次加载下一个页面时需要批量校验元组可见性之外,在后面每一次返回该页面下一条元组时,都要重新对页面加共享锁,判断下一条元组的可见性。该模式的查询性能较批量模式要稍低,适用于对系统表的顺序扫描(系统表的可见性不参照查询快照,而是以实时的事务提交状态为准)。 (4) 调用tableam_scan_end接口,从g_tableam_routines数组中找到astore的扫描结束接口,即heap_endscan接口,结束顺序扫描操作,释放对应的扫描结构体。 (5) 调用heap_close接口,释放对表加的锁或引用计数。 5. astore空间管理和回收 openGauss中采用最大堆二叉树结构来记录和管理astore堆表页面的空闲空间,该最大堆二叉树结构按照页面粒度进行与存储介质的读写操作,并单独储存于专门的空闲空间位图文件中(free space map,简称FSM)。该FSM文件的结构如图4-8所示。 图4-8 astore FSM文件结构示意图 所有页面分为叶子节点页面和内部节点页面两种。两种页面的页面内部结构完全相同,区别在于:对于叶子节点页面,其页面中记录的二叉树的叶子节点对应堆/索引表页面的空闲空间程度;对于内部节点页面,其页面中记录的二叉树的叶子节点对应下层FSM页面的最大空闲空间程度。 使用FSM页面中的1个字节(即256档)来记录一个堆/索引页面的空闲空间程度。在FSM页面中不会记录任何堆/索引页面的页号信息,也不会记录任何根、子FSM节点页面的页号信息,这些信息主要通过以下的规则来计算得到: (1) 在一个FSM页面内部,二叉树节点按照从上到下、从左到右逐层排布,即:第一个字节为根节点的空闲程度,第二个字节为第一层内部节点最左边节点的空闲程度,依次类推。 (2) 所有FSM页面在物理存储上采用深度优先顺序,即某个FSM页面之前所有的物理页面包括:该FSM页面所在子树的所有上层节点,加上该FSM页面所有左侧子树。 (3)所有FSM叶子节点页面中的二叉树的叶子节点,对应堆/索引表页面的空闲空间程度,且根据从左到右的顺序,分别对应第1个、第2个、….、第n个堆/索引表物理页面。 (4)除了(3)中这些FSM节点之外,其他FSM父节点保存子节点(子树)中空闲空间的最大值。 根据上述算法,可以高效地查询出具有足够空闲空间的堆/索引页面的页面号,并将待插入的数据插入其中。 FSM模块主要的对外接口和含义如表4-14所示。 表4-14 FSM模块主要的对外接口 接口名称 接口含义 GetPageWithFreeSpace 获取空闲程度大于入参的堆/索引页面号 RecordAndGetPageWithFreeSpace 更新当前(不满足条件的)堆/索引页面的空闲空间程度,寻找新的空闲程度大于入参的堆/索引页面号 RecordPageWithFreeSpace 更新单个堆/索引页面的空闲空间程度 UpdateFreeSpaceMap 更新多个(批量插入的)堆/索引页面的空闲空间程度 FreeSpaceMapTruncateRel 删除所有储存大于某个堆/索引页面号空闲信息的FSM页面 FreeSpaceMapVacuum 修正所有FSM内部节点的空闲空间信息 此外,为了保证FSM信息的维护操作不会带来明显的开销,因此FSM的所有修改都是不记录日志的。同时,对于某个堆/索引页面对应的FSM信息,只在页面初始化和页面空闲空间整理(见本节后面介绍)两种场景下才会主动更新,除此之外,只有当新插入的数据发现该页面实际空间不足时才会被动更新该页面对应的FSM信息(也包括由于宕机导致的FSM页面损坏)。 空闲空间的管理难点在于空闲空间的回收。在openGauss中,对于astore存储格式,有3种回收空闲空间的方式,如图4-9所示。 图4-9 astore空闲空间回收机制示意图 1. 轻量级堆页面清理 当查询扫描到某个astore堆表页面时,会顺带尝试清理该页面上已经被删除的、足够老的元组(足够老是指在元组对于所有并发查询均为已经删除状态,具体参见事务处理章节)。由于只是顺带清理该页面内容,因此只能删除元组内容本身,元组指针还需要保留,以免在索引侧造成空引用或空指针(可参见4.2.5 行存储索引机制)。一个比较特殊的情况是HOT场景。HOT场景是指对于该表上所有的索引更新前后的索引键值均没有发生变化,因此对于更新后的元组只需要插入堆表元组而不需要新插入索引元组。对于同一个页面内一条HOT链上的多个元组,如果它们都足够老了,那么在清理时可以额外删除所有中间的元组指针,只保留第一个版本的元组指针,并将其重定向到第一个不用被清理的元组版本的元组指针。 轻量级堆页面清理的接口是heap_page_prune_opt函数,关键的数据结构是PruneState结构体,定义代码如下: typedef struct { TransactionId new_prune_xid; TransactionId latestRemovedXid; int nredirected; /* 待重定向的元组个数 */ int ndead; /* 待标记死亡的元组个数 */ int nunused; /* 待回收的元组个数 */ OffsetNumber redirected[MaxHeapTuplesPerPage * 2]; OffsetNumber nowdead[MaxHeapTuplesPerPage]; OffsetNumber nowunused[MaxHeapTuplesPerPage]; bool marked[MaxHeapTuplesPerPage + 1]; } PruneState; 其中,“new_prune_xid”字段用于记录页面上此次没有被清理的、但是已经被删除的元组的xmax,用于决定下次何时再次清理该页面;“latestRemovedXid”字段用于记录该页面上被清理的元组的最大的xmax,用于判断热备上回放页面整理时是否需要等待只读查询;nredirected、ndead、nunused、redirected、nowdead和nowunused分别记录该页面上待重定向的、待标记死亡的、待回收的元组。 2. 中量级堆页面和索引页面清理 openGauss提供VACUUM语句来让用户主动执行对某个astore表(或某个库中所有的astore表)及其上的索引进行中量级清理。中量级清理过程中,不阻塞相关表的查询和DML操作。由于在astore表中,新、老版本元组是混合存储的,因此,与顺带执行的轻量级清理相比,astore表的中量级清理需要进行全表顺序(或索引)扫描,才能识别出所有待清理的老版本元组。对于扫描出来的确认要清理的元组,会首先清理索引中的元组,然后再清理堆表中的元组,从而可以避免出现索引空指针的问题。 中量级清理的对外接口是lazy_vacuum_rel函数,内部逐层调用lazy_scan_rel、lazy_scan_heap和heap_page_prune(同轻量级清理)来扫描和暂存几类待清理的元组。当待清理的元组积攒到一定数量之后(受maintenance_work_mem内存上限控制),先后调用lazy_vacuum_index接口和lazy_vacuum_heap接口来分别清理索引文件和堆表文件。其中,与堆表页面将元组指针置为UNUSED不同,索引页面直接删除被清理的元组指针,并进行页面重整。 中量级清理的关键数据结构是LVRelStats结构体,定义代码如下: typedef struct LVRelStats { bool hasindex; /* 表上是否有索引 */ /* 统计信息 */ BlockNumber old_rel_pages; /* 之前的页面个数统计 */ BlockNumber rel_pages; /* 当前的页面个数统计 */ BlockNumber scanned_pages; /* 已经扫描的页面个数 */ double scanned_tuples; /* 已经扫描的元组个数 */ double old_rel_tuples; /* 之前的元组个数统计 */ double new_rel_tuples; /* 当前的元组个数统计 */ BlockNumber pages_removed; double tuples_deleted; BlockNumber nonempty_pages; /* 最后一个非空页面的页面号加1 */ /* 待清理的元组的行号信息(已排序) */ int num_dead_tuples; /* 当前待清理的元组个数 */ int max_dead_tuples; /* 单次最多可记录的待清理元组个数 */ ItemPointer dead_tuples; /* 待清理元组行号数组 */ int num_index_scans; TransactionId latestRemovedXid; bool lock_waiter_detected; BlockNumber* new_idx_pages; double* new_idx_tuples; bool* idx_estimated; Oid currVacuumPartOid; } LVRelStats; 其中hasindex表示该表是否有索引表,num_dead_tuples表示目前已经积攒的要清理的元组,dead_tuples是保存这些元组位置的TID数组,max_dead_tuples是根据maintenance_work_mem计算出来的单次允许积攒的最大待清理元组个数。 需要指出的是,如果在元组更新时就把老版本元组集中存储,那么清理时就无须全表扫描,只需要清理集中存储的老版本元组页面即可,这样可以有效降低清理过程带来的I/O开销,使得整体存储引擎的I/O开销和性能更平稳,这也是后续openGauss版本将支持的ustore行存储格式的设计出发点。 3. 重量级堆页面和索引页面清理 无论是轻量级清理,或是中量级清理,都只能局部清理astore页面中的死亡元组,无法真正实现对这些空闲空间的释放(被清理出的空间,仍然只能被该表使用)。因此,openGauss还提供了VACUUM FULL语句来让用户主动执行对某个astore表(或某个库中所有astore表)及其上的索引进行重量级清理。重量级清理将一个表中所有仍未死亡(但是可能已经被删除)的元组重新紧密插入到新的堆表文件中并在此基础上重新创建所有索引,从而实现对空闲空间的彻底回收。在重量级清理的主体流程中只允许用户执行只读查询操作,在重量级清理的提交流程中只读查询操作也会被阻塞。 为了尽可能提高重新创建的索引性能,如果用户堆表上有索引,那么上述全表扫描会采用索引扫描。 重量级清理的对外接口是cluster_rel函数,内部逐层调用rebuild_relation、copy_heap_data、tableam_relation_copy_for_cluster、heapam_relation_copy_for_cluster、copy_heap_data_internal、reform_and_rewrite_tuple、rewrite_heap_tuple。其中,“rewrite_heap_tuple”接口将每一条扫描的未死亡元组进行重构(去除被删除的字段)之后,插入到新的紧密排列的堆表中。在这个过程中,对原来多个元组之间的更新链关系采用两个哈希表来进行暂存。当一对更新元组的双方都扫描到之后,就进行新表的填充,并将更新后元组的新的TID(transaction ID,事务ID)保存到更新前的元组中。上述机制保证重量级清理过程中并发更新事务的执行机制不会受到破坏。 重量级清理的关键数据结构是RewriteStateData结构体,其定义代码如下: typedef struct RewriteStateData { Relation rs_old_rel; /* 源表 */ Relation rs_new_rel; /* 整理后的目标表 */ Page rs_buffer; /* 当前整理的源表页面 */ BlockNumber rs_blockno; /* 当前写入的目标表页面号 */ bool rs_buffer_valid; /* 当前缓冲区是否有效 */ bool rs_use_wal; /* 整理操作是否产生日志 */ TransactionId rs_oldest_xmin; /* 用于可见性判断的最老活跃事务号 */ TransactionId rs_freeze_xid; /* 用于元组冻结判断的事务号 */ MemoryContext rs_cxt; /* 哈希表内存上下文 */ HTAB *rs_unresolved_tups; /* 未匹配的更新前元组版本 */ HTAB *rs_old_new_tid_map; /* 未匹配的更新后元组版本 */ /* 元组压缩相关信息 */ PageCompress *rs_compressor; Page rs_cmprBuffer; HeapTuple *rs_tupBuf; Size rs_size; int rs_nTups; bool rs_doCmprFlag; /* 异步-同步读写相关 */ char *rs_buffers_queue; /* adio write queue */ char *rs_buffers_queue_ptr; /* adio write queue ptr */ BufferDesc *rs_buffers_handler; /* adio write buffer handler */ BufferDesc *rs_buffers_handler_ptr; /* adio write buffer handler ptr */ int rs_block_start; /* adio write start block id */ int rs_block_count; /* adio write block count */ } RewriteStateData; 其中,rs_old_rel是被清理的表,rs_new_rel是清理之后的表,rs_oldest_xmin是判断元组是否死亡的xid阈值,rs_freeze_xid是判断是否进行freeze操作的xid阈值。rs_unresolved_tups是保存一对更新元组中老元组的哈希表,rs_old_new_tid_map是保存一对更新元组中新元组的哈希表,这两个成员共同保证更新链信息不被丢失(在原表中更新后的元组的物理位置可能比更新前的元组的物理位置还要小)。 最后,重量级操作实际上是一种数据重聚簇操作,对于其他行存储子格式和cstore列存储格式同样适用,只是具体实现机制略有不同。 本期精彩内容将告一段落,下篇我们将详细介绍“4.2.4 ustore”相关内容,敬请期待!

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

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元组多版本机制”相关内容,敬请期待!

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

【Java技术探索】ThreadPoolExecutor深入浅出的源码分析(核心源码)

线程池执行任务的流程 如果线程池工作线程数<corePoolSize,创建新线程执行task,并不断轮训t等待队列处理task。 如果线程池工作线程数>=corePoolSize并且等待队列未满,将task插入等待队列。 如果线程池工作流程数>=corePoolSize并且等待队列已满,且工作线程数<maximumPoolSize,创建新线程执行task。 如果线程池工作流程数>=corePoolSize并且等待队列已满,且工作线程数=maximumPoolSize,执行拒绝策略。 execute()原理 public void execute(Runnable command) { if (command == null) throw new NullPointerException(); /* * Proceed in 3 steps: * 1. If fewer than corePoolSize threads are running, try to * start a new thread with the given command as its first * task.The call to addWorker atomically checks runState and * workerCount, and so prevents false alarms that would add * threads when it shouldn't, by returning false. * * 如果运行的线程数小于corePoolSize,尝试创建一个新线程(Worker),并执行 * 它的第一个任务command。 * * 2. If a task can be successfully queued, then we still need * to double-check whether we should have added a thread * (because existing ones died since last checking) or that * the pool shut down since entry into this method. So we * recheck state and if necessary roll back the enqueuing if * stopped, or start a new thread if there are none. * * 如果task成功插入等待队列,我们仍需要进行双重校验是否可以成功添加一个线程 * (因为有的线程可能在我们上次检查以后已经死掉了)或者在我们进入这个 * 方法后线程池已经关闭了。 * * 3. If we cannot queue task, then we try to add a new * thread. If it fails, we know we are shut down or saturated * and so reject the task. * 如果等待队列已满,我们尝试新创建一个线程。如果创建失败,我们知道线程已关闭 或者已饱和,因此我们拒绝改任务。 */ int c = ctl.get(); // 高3位表示状态,低29位任务数量。 //工作线程小于核心线程数,创建新的线程。 if (workerCountOf(c) < corePoolSize) { //创建新的worker立即执行command,并且轮训workQueue处理task。 (核心线程数) if (addWorker(command, true)) return; c = ctl.get(); } //线程池在运行状态且可以将task插入队列 //第一次校验线程池在运行状态 if (isRunning(c) && workQueue.offer(command)) { int recheck = ctl.get(); //第二次校验,防止在第一次校验通过后线程池关闭。 如果线程池关闭,在队列中删除task并拒绝task if (! isRunning(recheck) && remove(command)) reject(command); //如果线程数=0(线程都死掉了,比如:corePoolSize=0),新建线程且未指 // 定firstTask,仅仅去轮训workQueue else if (workerCountOf(recheck) == 0) addWorker(null, false); } //线程队列已满,尝试创建新线程执行task,创建失败后拒绝task //创建失败原因:1.线程池关闭;2.线程数已经达到maxPoolSize else if (!addWorker(command, false)) reject(command); } addWorker()原理 线程管理部分的分水岭-方法调用之前 都是任务管理(任务的创建、以及拒绝、添加任务队列等,任务线程池的状态监控等),方法调用之后属于线程的管理。(线程的执行和阻塞等) 参数: firstTask:worker线程的初始任务,可以为空。 core:true:将corePoolSize作为上限,false:将maximumPoolSize作为上限 private boolean addWorker(Runnable firstTask, boolean core) { retry: //外层循环判断线程池的状态 for (;;) { //在进行获取一次上下文操作机制 int c = ctl.get(); int rs = runStateOf(c); //线程池状态 // Check if queue empty only if necessary. // 线程池状态:RUNNING = -1、SHUTDOWN = 0、 STOP = 1、TIDYING = 2、TERMINATED = 3 //线程池至少是shutdown状态 if (rs >= SHUTDOWN && // 除了线程池正在关闭(shutdown), // 队列里还有未处理的task的情况,其他都不能添加 ! (rs == SHUTDOWN && firstTask == null && ! workQueue.isEmpty())) return false; //内层循环判断是否到达容量上限,worker+1 for (;;) { int wc = workerCountOf(c);//worker数量 //worker大于Integer最大上限 //或到达边界上限 if (wc >= CAPACITY || wc >= (core ? corePoolSize : maximumPoolSize)) return false; //CAS worker+1 if (compareAndIncrementWorkerCount(c)) break retry;//成功了跳出循环 c = ctl.get(); // Re-read ctl if (runStateOf(c) != rs) //如果线程池状态发生变化,重试外层循环 continue retry; // else CAS failed due to workerCount change; // retry inner loop // CAS失败workerCount被其他线程改变, // 重新尝试内层循环CAS对workerCount+1 } } boolean workerStarted = false; boolean workerAdded = false; Worker w = null; try { final ReentrantLock mainLock = this.mainLock; w = new Worker(firstTask); //1.state置为-1,Worker继承了AbstractQueuedSynchronizer. //2.设置firstTask属性. //3.Worker实现了Runable接口,将this作为入参创建线程. final Thread t = w.thread; if (t != null) { //addWorker需要加锁 mainLock.lock(); try { // Recheck while holding lock. // Back out on ThreadFactory failure or if // shut down before lock acquired. int c = ctl.get(); int rs = runStateOf(c); if (rs < SHUTDOWN || (rs == SHUTDOWN && firstTask == null)) { if (t.isAlive()) // precheck that t is startable throw new IllegalThreadStateException(); workers.add(w);//workers是HashSet<Worker> //设置最大线程池大小 int s = workers.size(); if (s > largestPoolSize) largestPoolSize = s; workerAdded = true; } } finally { mainLock.unlock(); } if (workerAdded) { t.start(); workerStarted = true; } } } finally { if (!workerStarted) addWorkerFailed(w); } return workerStarted; } 可以理解就是每一个Worker对象都是一个AQS队列哦! addWorker方法有4种传参的方式: addWorker(command, true) 线程数小于corePoolSize。判断workers(HashSet<Worker>)大小,如果worker数量>=corePoolSize返回false,否则创建worker添加到workers,并执行worker的run方法(执行firstTask并轮询tworkQueue); addWorker(command, false) 线程数大于corePoolSize且workQueue已满。如果worker数量>=maximumPoolSize返回false,否则创建worker添加到workers,并执行worker的run方法(执行firstTask并轮询tworkQueue); addWorker(null, false) 没有worker存活也就是任务梳理runcount为0,创建worker去轮询workQueue,长度限制maximumPoolSize。 addWorker(null, true) 在execute方法中就使用了前3种,结合这个核心方法进行以下分析 以上无论哪种方式都需要进行相关的ReentrantLock的加锁,所以效率和性能不会特别好。所以有了一个小办法,prestartAllCoreThreads() prestartAllCoreThreads()原理 这个方法调用,启动所有的核心线程去轮询workQueue。因为addWorker是需要上锁的,预启动核心线程可以提高执行效率。 ThreadPoolExecutor 内部类Worker (线程管理的核心类) /** * Class Worker mainly maintains interrupt control state for * threads running tasks, along with other minor bookkeeping. * This class opportunistically extends AbstractQueuedSynchronizer * to simplify acquiring and releasing a lock surrounding each * task execution. This protects against interrupts that are * intended to wake up a worker thread waiting for a task from * instead interrupting a task being run. We implement a simple * non-reentrant mutual exclusion lock rather than use * ReentrantLock because we do not want worker tasks to be able to * reacquire the lock when they invoke pool control methods like * setCorePoolSize. Additionally, to suppress interrupts until * the thread actually starts running tasks, we initialize lock * state to a negative value, and clear it upon start (in * runWorker). * 1.Worker类主要负责运行线程状态的控制。 * 2.Worker继承了AQS实现了简单的获取锁和释放所的操作。来避免中断等待执行任务的线 * 程时,中断正在运行中的线程(线程刚启动,还没开始执行任务)。 * 3.自己实现不可重入锁,是为了避免在实现线程池控状态控制的方法,例如: * setCorePoolSize的时候中断正在开始运行的线程。 * setCorePoolSize可能会调用interruptIdleWorkers(),该方法中会调用worker的tryLock()方法 * 中断线程,自己实现锁可以确保工作线程启动之前不会被中断 */ private final class Worker extends AbstractQueuedSynchronizer implements Runnable { /** * This class will never be serialized, but we provide a * serialVersionUID to suppress a javac warning. */ private static final long serialVersionUID = 6138294804551838833L; /** Thread this worker is running in. Null if factory fails. */ 封装任务线程机制。 final Thread thread; /** Initial task to run. Possibly null. */ Runnable firstTask; /** Per-thread task counter */ volatile long completedTasks; /** * Creates with given first task and thread from ThreadFactory. * @param firstTask the first task (null if none) */ Worker(Runnable firstTask) { // inhibit interrupts until runWorker //状态置为-1,如果中断线程需要CAS将state 从0- >1,以此来保证能只中断从workerQueue getTask的线程 setState(-1); this.firstTask = firstTask; this.thread = getThreadFactory().newThread(this); } // 核心方法机制 /** Delegates main run loop to outer runWorker */ public void run() { //首先执行w.unlock,就是把state置为0,对该线程的中断就可以进行了 runWorker(this); } // Lock methods // // The value 0 represents the unlocked state. // The value 1 represents the locked state. protected boolean isHeldExclusively() { return getState() != 0; } // 在setCorePoolSize/shutdown等方法中断worker线程时需要调用该方法, // 确保中断的是从workerQueue getTask的线程 protected boolean tryAcquire(int unused) { if (compareAndSetState(0, 1)) { setExclusiveOwnerThread(Thread.currentThread()); return true; } return false; } protected boolean tryRelease(int unused) { setExclusiveOwnerThread(null); setState(0); return true; } public void lock() { acquire(1); } public boolean tryLock() { return tryAcquire(1); } public void unlock() { release(1); } //调用tryRelease修改state=0,LockSupport.unpark(thread) 下一个等待锁的线程 public boolean isLocked() { return isHeldExclusively(); } void interruptIfStarted() { Thread t; if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) { try { t.interrupt(); } catch (SecurityException ignore) { } } } } 阻塞队列 workQueue 有多种选择,在 JDK 中一共提供了 7 中阻塞对列,分别为: ArrayBlockingQueue : 一个由数组结构组成的有界阻塞队列。 此队列按照先进先出(FIFO)的原则对元素进行排序。默认情况下不保证访问者公平地访问队列 ,所谓公平访问队列是指阻塞的线程,可按照阻塞的先后顺序访问队列。非公平性是对先等待的线程是不公平的,当队列可用时,阻塞的线程都可以竞争访问队列的资格。 LinkedBlockingQueue : 一个由链表结构组成的有界阻塞队列。 此队列的默认和最大长度为Integer.MAX_VALUE。 此队列按照先进先出的原则对元素进行排序。 PriorityBlockingQueue : 一个支持优先级排序的***阻塞队列。 (虽然此队列逻辑上是***的,但是资源被耗尽时试图执行 add 操作也将失败,导致 OutOfMemoryError) DelayQueue: 一个使用优先级队列实现的***阻塞队列。 元素的一个***阻塞队列,只有在延迟期满时才能从中提取元素 SynchronousQueue: 一个不存储元素的阻塞队列。 一种阻塞队列,其中每个插入操作必须等待另一个线程的对应移除操作 ,反之亦然。(SynchronousQueue 该队列不保存元素) LinkedTransferQueue: 一个由链表结构组成的***阻塞队列。 相对于其他阻塞队列LinkedTransferQueue多了tryTransfer和transfer方法。 LinkedBlockingDeque: 一个由链表结构组成的双向阻塞队列。 是一个由链表结构组成的双向阻塞队列 关闭线程池 其实,如果优雅的关闭线程池是一个令人头疼的问题,线程开启是简单的,但是想要停止却不是那么容易的。通常而言, 大部分程序员都是使用 jdk 提供的两个方法来关闭线程池,他们分别是:shutdown 或 shutdownNow; 通过调用线程池的 shutdown 或 shutdownNow 方法来关闭线程池。它们的原理是遍历线程池中的工作线程,然后逐个调用线程的 interrupt 方法来中断线程(PS:中断,仅仅是给线程打上一个标记,并不是代表这个线程停止了,如果线程不响应中断,那么这个标记将毫无作用),所以无法响应中断的任务可能永远无法终止。但是它们存在一定的区别,shutdownNow首先将线程池的状态设置成 STOP,然后尝试停止所有的正在执行或暂停任务的线程,并返回等待执行任务的列表,而 shutdown 只是将线程池的状态设置成SHUTDOWN状态,然后中断所有没有正在执行任务的线程。只要调用了这两个关闭方法中的任意一个,isShutdown 方法就会返回 true。当所有的任务都已关闭后,才表示线程池关闭成功,这时调用isTerminaed方法会返回 true。至于应该调用哪一种方法来关闭线程池,应该由提交到线程池的任务特性决定,通常调用 shutdown方法来关闭线程池,如果任务不一定要执行完,则可以调用 shutdownNow 方法。这里推荐使用稳妥的 shutdownNow 来关闭线程池,至于更优雅的方式我会在以后的并发编程设计模式中的两阶段终止模式中会再次详细介绍。 线程池的优点 在 Java 并发编程框架中的线程池是运用场景最多的技术,几乎所有需要异步或并发执行任务的程序都可以使用线程池。在开发过程中,合理地使用线程池能够带来至少以下4个好处。 第一:降低资源消耗。通过重复利用已创建的线程降低线程创建和销毁造成的消耗; 第二:提高响应速度。当任务到达时,任务可以不需要等到线程创建就能立即执行; 第三:提高线程的可管理性。线程是稀缺资源,如果无限制地创建,不仅会消耗系统资源,还会降低系统的稳定性,使用线程池可以进行统一分配、调优和监控。 第四:提供更强大的功能,比如延时定时线程池; 线程池的场景特性 任务的性质:CPU密集型任务、IO密集型任务和混合型任务; 任务的优先级:高、中和低; 任务的执行时间:长、中和短; 任务的依赖性:是否依赖其他系统资源,如数据库连接; 性质不同的任务可以用不同规模的线程池分开处理。分为CPU密集型和IO密集型。 CPU密集型任务应配置尽可能小的线程,如配置 Ncpu+1个线程的线程池。(可以通过Runtime.getRuntime().availableProcessors()来获取CPU物理核数),参考建议哈 IO密集型任务线程并不是一直在执行任务,则应配置尽可能多的线程,如 2*Ncpu。 混合型的任务,如果可以拆分,将其拆分成一个CPU密集型任务一个IO密集型任务,只要这两个任务执行的时间相差不是太大,那么分解后执行的吞吐量将高于串行执行的吞吐量。 如果这两个任务执行时间相差太大,则没必要进行分解。 可以通过 Runtime.getRuntime().availableProcessors() 方法获得当前设备的CPU个数。 优先级不同的任务可以使用优先级队列 PriorityBlockingQueue来处理。它可以让优先级高的任务先执行(注意:如果一直有优先级高的任务提交到队列里,那么优先级低的任务可能永远不能执行)执行时间不同的任务可以交给不同规模的线程池来处理,或者可以使用优先级队列,让执行时间短的任务先执行。 依赖数据库连接池的任务,因为线程提交SQL后需要等待数据库返回结果,等待的时间越长,则 CPU 空闲时间就越长,那么线程数应该设置得越大,这样才能更好地利用CPU。 建议使用有界队列。有界队列能增加系统的稳定性和预警能力,可以根据需要设大一点。方式因为提交的任务过多而导致 OOM;

资源下载

更多资源
腾讯云软件源

腾讯云软件源

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

Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

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等操作系统。

用户登录
用户注册