0%

阿土登录流计算平台后,看到平台上面可以编写 Sql 语法,于是就写了一个简单的 sql。

图片

他发现旁边有个校验功能这时平台弹出 SQL 语法效验正确。

SQL 语法效验完成后,阿土点击提交按钮,流计算平台提示,SQl 语法效验正确,已成功提交集群。

Flink sql 代码居然提交到yarn集群上了???

图片

  1. Flink Sql解析器

  2. Flink Planner 和 Blink Planner

  3. Blink Sql提交流程

1.1、了解Calcite

为方便用户使用 Flink 流计算组件,Flink 社区设计了四种抽象,在这些抽象中,Sql API 属于Flink的最上层抽象,是 Flink 的一等公民,这就方便用户或者开发者直接通过 Sql 编写来提交任务。

图片

但经过阿土的调查后 发现,Flink sql 在提交任务时,并不是向 DataStream API 那样,直接被转为 StreamGraph,经过优化生成 JobGraph 提交到集群的,而是需要对编写的 Sql 进行解析、验证、优化等操作,在这中间,社区引入了一个强大的解析器,那就是Calcite

1.2、Calcite执行步骤

Calcite的执行流程,主要涉及5个部分 SQL解析、SQL校验、SQL查询优化、SQL生成、执行等。

图片

在这个流程中,Calcite各阶段扮演的角色如下:

  1. SQL解析。通过 JavaCC 实现,使用 JavaCC 编写 SQL 语法描述文件,将 SQL 解析成未经校验的 AST 语法树

  2. SQL校验。通过与元数据结合验证 SQL 中的 Schema、Field、 Function 是否存在,输入输出类型是否匹配等。

  3. SQL优化。对上个步骤的输出( RelNode ,逻辑计划树)进行优化,使用两种规则:基于规则优化 和 基于代价优化,得到优化后的物理执行计划。

  4. SQL生成。将物理执行计划生成为在特定平台/引擎的可执行程序,如生成符合 MySQL 或 Oracle 等不同平台规则的 SQL 查询语句等。

  5. 执行。执行是通过各个执行平台执行查询,得到输出结果。

其中,Calcite再与其他处理引擎结合时,到SQL优化阶段就已经结束。所以流程图简化为:

图片

Calcite是怎么在Flink中扮演的角色呢?Flink1.13.2Flink Sql`源码

2.1 Flink Planner和Blink Planner

在1.9.0版本以前,社区使用Flink Planner作为查询处理器,通过与Calcite进行连接,为Table/SQL API提供完整的解析、优化和执行环境,使其SQL被转为DataStream API的 Transformation,然后再经过StreamJraph -> JobGraph -> ExecutionGraph等一系列流程,最终被提交到集群。

在1.9.0版本,社区引入阿里巴巴的Blink,对FIink TabIe & SQL模块做了重大的重构,保留了 Flink Planner 的同时,引入了 Blink Planner,没引入以前,Flink 没考虑流批作业统一,针对流批作业,底层实现两套代码,引入后,基于流批一体理念,重新设计算子,以流为核心,流作业和批作业最终都会被转为transformation

2.2 Blink Planner与Calcite关系

在之后的版本,为了实现Flink流批一体的愿景,通过Blink Planner与Calcite进行对接,对接流程如下:

  1. 在Table/SQL 编写完成后,通过Calcite 中的parse、validate、rel阶段,以及Blink额外添加的convert阶段,将其先转为Operation

  2. 通过Blink Planner 的translateToReloptimizetranslateToExecNodeGraphtranslateToPlan四个阶段,将Operation转换成DataStream API的 Transformation

  3. 再经过StreamGraph -> JobGraph -> ExecutionGraph等一系列流程,SQL最终被提交到集群。

SQL执行流程图

图片

根据对源码的分析后,发现无论是Flink SQL执行DDL操作、还是DQL操作或者DML操作、最终都可以将其总结为两个阶段:

  1. SQL 语句到 Operation 过程,即Parse阶段;

  2. Operation 到 Transformations 过程,即Translate阶段。

3.1、Parse阶段

在Parse阶段一共包含parse、validate、rel、convert部分

图片

Calcite的 parse 解析模块是基于javacc实现的。javacc是一个词法分析生成器语法分析生成器。词法分析器于将输入字符流解析成一个一个的token,以下面这段SQL语句为例:

示例1 :

图片

在 parse 部分,上面的SQL语句最后会被解析为如下一组token:

图片

接下来语法分析器会以词法分析器解析出来的token序列作为输入来进行语法分析。分析过程使用递归下降语法解析,LL(k)。

其中,第一个L表示从左到右扫描输入;第二个L表示每次都进行最左推导(在推导语法树的过程中每次都替换句型中最左的非终结符为终结符。类似还有最右推导);

k表示的是每次向前探索(lookahead)k 个终结符。

分析所依赖的的词法法则定义在一个parser.jj文件中。

图片

在经过词法分析和语法分析后,一段 SQL 语句会被解析成一颗抽象语法树(Abstract Syntax Tree,AST),树的节点类型在 Calcite 中以 SqlNode 来表示,不同节点以不同子类型的SqlNode来表示。

同样以上面的SQL为例,在这段SQL中:

  1. id, score, T 等为 SqlIdentifier,表示一个字段名或表名的标识符;

  2. select和cast()为SqlCall,表示一个行为或动作,其中cast()为一个SqlBasicCall,表示一个函数调用,具体调用的是什么函数,由其内部的SqlOperator决定,比如这里是一个二元操作符“<”,对应SqlBinaryOperator,operator的名字是“<”,类别是SqlKind.LESS_THAN;

  3. int 为 SqlDataTypeSpec,表示一个类型定义;

  4. ‘hello’和 10 为SqlLiteral,表示一个常量;

在Calcite中,所有的操作都是一个SqlCall, 如查询是一个 SqlSelect, 删除是一个 SqlDelete 等,它们都是 SqlCall 的子类型。select的查询条件等为 SqlCall 中的参数。示例1 的 SQL 语句最终生成的语法树形式如下:

图片

如果把示例1中的直接从一个表查询数据,改为从两张表的关联结果中查询数据,例如:

示例2:

图片

则相应的AST形式如下:

图片

其中只有FROM子树部分由原来的SqlIdentifier节点变成了一棵SqlJoin子树,其他部分与示例1相同所以在图中省略了。

校验(validate)阶段

图片

对经过parser解析出的AST进行有效性验证,验证的方面主要包括以下两方面:

  1. 表名、字段名、函数名是否正确,如在某个查询的字段在当前SQL位置上是否存在或有歧义(当前可见的多个数据源中同时存在该名称的字段)

  2. 特定类型操作自身的合法性,如group by聚合中的聚合函数是否存在嵌套调用,使用AS重命名时,新名字是否是x.y的形式等

针对上面的第一种情况,在校验过程中首先需要明确两个最重要的概念:NameSpaceScope

NameSpace代表一个逻辑上的数据源,可以是一张表,也可以是一个子查询,而Scope则代表了在 SQL 的某个位置,表和字段的可见范围。

从概念中可以看出,在某个 SQL位置上,某个字段所对应的 scope 可能包含多个 namespace。在 validate 阶段解析出来的 scope 和 namespace 信息会被保存下来,在后面转换成逻辑执行计划的时候还会用到。

通过一个示例来看什么是 NameSpace 和 scope

示例3

图片

在上面这样一段SQL语句中包含四个namespace:

图片

对于SQL中的不同表达式,根据它们所在的位置,它们所对应的scope如下:

图片

那么在校验第一种情况的时候,整个校验过程的核心就在于为不同的SqlNode节点生成其对应的namespace和scope,然后对该SqlNode涉及的字段和namespacescope的对应关系进行校验。

对于第二种情况的校验,则需要根据具体的节点类型分别实现了。

在Calcite中,validator的具体实现类是SqlValidatorImpl,namespace和scope分别由接口SqlValidatorNamespaceSqlValidatorScope表示,图中涉及到的xxxNamespace和xxxScope分别是这两个类的子类。

下图是从调用validator.validate(sqlNode)开始,对一段查询语句的表名和字段名进行校验的时序图。

图片

大体过程都已经在图中的注解里进行了说明,需要补充的一点是,在通过 emptyScope.resolve解析表名时,表信息是通过具体的catalogReader从catalog的schema中查找出来的。

具体使用什么catalog和catalogReader,是在validator创建之初决定的

在flink中,根据用户的配置,catalog可能是 GenericInMemoryCatalog(基于内存的catalog)或HiveCatalog(基于hive metastore的catalog)。

如下图所示:

图片

rel阶段是将SqlNode组成的一棵抽象语法树转化为一棵由RelNode和RexNode组成的关系代数树,或者称为执行计划。RelNode表示关系表达式,如投影(Project),即SELECT,和连接(JOIN)等;

RexNode表示行表达式,如示例中的 CAST(score AS INT)、T1.id < 10

以示例2的语法树为例,在经过rel阶段转换后会生成下图所示的执行计划:

图片

rel阶段只处理DML和DQL

因为DDL实际上可以认为是对元数据的修改,不涉及复杂关系查询,也就不用进行关系代数转换来优化执行,所以也无需转换为表示,根据对应的SqlNode中保存的信息已经可以直接执行了。

在calcite中,SqlToRelConverter用于对关系表达式进行转换。Flink中通过如下方式使用calcite将AST转换成逻辑执行计划,如下图源码所示。

图片

从Flink1.13.2源码中可以看到转换的入口是convertQuery方法。

SqlToRelConverter中的简单的转换流程如下图所示:

图片

针对每种可能的根节点类型都有对应的转换方法。其中DELETEUPDATEMERGEWITHVALUES这几种语法在flink流式SQL中还不支持,并且其转换过程也比较简单,后文不再详细分析。

对于一棵转换后得到的逻辑执行计划树中的节点,其实在AST中都是可以一一对应的找到对应的节点的,所以转换过程本身并不涉及很复杂的算法,大部分过程是提取已有SqlNode节点中记录的信息,然后生成对应的RelNodeRexNode,并设置RelNode间的父子关系。

从图中也可以看出在calcite里最终都会生成一个LogicalModify节点,通过节点内的operation属性来标识不同的含义。但是目前flink支持的DML只有insert语句,而且并不会生成LogicalModify节点,而是直接转换成了ModifyOperation,并在需要的时候转换成flink内部自己定义的节点类型LogicalSink。也因为这个原因,对于DML的转换流程图中是略有简化的,insertdeleteupdatemerge本身都可以带查询语句,因此实际转换的时候都会递归地先对查询部分进行转换。

上图所示流程中只展示了对关系表达式的转换,但是每个关系节点(RelNode)中的行表达式同样需要经过转换得来。

Calcite中行表达式的转换依赖于两个对象:BlackBoardSqlNodeToRexConverter

BlackBoard是对select进行转换时的一个临时工作空间,它就像一块“黑板”一样,可以临时记录下转换过程中需要的信息,比如select依赖的scope、当前的root节点、当前节点是否是top节点等。

BlackBoard本身还是一个shuttle,针对不同类型的SqlNode,其内部都有对应的visit方法。其中除SqlCallSqlLiteralSqlIntervalQualifier外,都可由BlackBoardSqlToRelConverter中定义的各种convertXXX方法进行转换,这三种类型的SqlNode则需要借助SqlNodeToRexConverter来进行转换。

SqlLiteralSqlIntervalQualifier的转换比较简单,就是从原来的SqlNode中提取信息进行简单的处理和转换,然后生成对应的RexNode

重点:这一步负责将RelNode tree转换成operation

图片

RelNode转换成Operation的过程很简单,针对四种类型的操作,其各自的转换过程如下:·

  1. CreateTable @convertCreateTable

    如果AST的根节点是SqlCreateTable,提取节点中记录的schemapropertiescommentprimary keysif not exists信息,创建CatalogTable对象,然后创建CreateTableOperation

  2. DropTable @convertDropTable

    如果AST的根节点是SqlDropTable,提取节点中记录的full table nameif exists信息,创建DropTableOperation对象

  3. Insert @convertInsert

    如果AST的根节点是RichSqlInsert,提取节点中记录的目标表的完整路径和查询表达式,先将查询表达式通过convertSqlQuery转换成QueryOperation,然后以转换后的QueryOperation为子节点创建ModifyOperation对象。

这里分两种情况:

(1)使用SQL API执行了 insert into 语句,将数据写入已经通过 TableEnvironment注册过的表中,此时创建的是CatalogSinkModifyOperation

(2)使用Table API的toXXXStream将table对象转换成了DataStream,创建的是OutputConversionModifyOperation

  1. Query @convertSqlQuery

    如果根节点的SqlKind是SqlKind.Query,先通过FlinkPlannerImpl.rel将SqlNode转换成RelNode,然后创建PlannerQueryOperation对象

3.2、Translate阶段

在Translate阶段,通过Blink Planner 的translateToReloptimizetranslateToExecNodeGraphtranslateToPlan四个阶段:将Operation转换成 Transformations。

重点

  1. 从operation开始,先将ModifyOperation通过translateToRel方法转换成Calcite RelNode逻辑计划树,在对应转换成FlinkLogicalRel(RelNode逻辑计划树);

  2. 然后经过 调用optimize方法将FlinkLogicalRel 优化成FlinkPhysicalRel。

  3. 再调用translateToExecNodeGraph方法将FlinkPhysicalRel转为execGraph

  4. 最后调用translateToPlan方法将execGraph转为transformations

图片

图片

从逻辑计划变成物理计划(RelNode),

图片

Flink1.13.2源码如下:

图片

这个过程可以看成是convert: RelNode => Operation的逆过程。

逻辑也很简单,无论是使用SQL API还是Table API,最终生成的operation的根节点一定是ModifyOperation,因为只有insert语句或者将Table转换成DataStream后,在DataStream结果上面写入sink才能触发执行。

前文提到过ModifyOperation最终都会被转换成flink内自定义的LogicalSink节点,该节点主要记录数据输出信息,核心在于需要创建出表示数据输出的sink。所以针对三种ModifyOperation类型分别创建sink的过程如下:

  1. UnregisteredSinkModifyOperation

    这个operation中直接记录了sink信息,因此直接提取出来创建LogicalSink即可。

  2. CatalogSinkModifyOperation

    根据operation中记录的table path找到对应的table,然后根据table创建出table sink,最后使用table sink创建出LogicalSink节点。

    这个过程中涉及到了在catalog中解析table和使用ServiceLoader根据table信息在classpath中查找并用于创建table sink的TableSinkFactory的过程,具体如下图所示。

图片

会使用两个优化器:RBO(基于规则的优化器) 和 CBO(基于代价的优化器)

  1. RBO(基于规则的优化器)会将原有表达式裁剪掉,遍历一系列规则(Rule),只要满足条件就转换,生成最终的执行计划。一些常见的规则包括分区裁剪(Partition Prune)、列裁剪、谓词下推(Predicate Pushdown)、投影下推(Projection Pushdown)、聚合下推、limit下推、sort下推、常量折叠(Constant Folding)、子查询内联转join等。

2.CBO(基于代价的优化器)会将原有表达式保留,基于统计信息和代价模型,尝试探索生成等价关系表达式,最终取代价最小的执行计划。CBO的实现有两种模型,Volcano模型,Cascades模型。这两种模型思想很是相似,不同点在于Cascades模型一边遍历SQL逻辑树,一边优化,从而进一步裁剪掉一些执行计划。

源码如下:

图片

调用translateToExecNodeGraph方法将FlinkPhysicalRel转为execGraph

图片

调用translateToPlan方法将execGraph转为transformations

图片

通过上述四个步骤,实现将Operation转换成 Transformations。

小笨猪通过完整的流程分析后,终于搞懂了Flink sql的解析和转换过程,最终SQL被转为Transformations,后面的步骤就变成了Flink DataStream的提交流程,小笨猪还是比较了解的。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
//
org.apache.flink.streaming.api.environment.StreamExplainEnvironment
//
try {
streamExplainEnvironment.setAsContext();

try {
// do something.
streamExplainEnvironment.execute("");
} catch (InvocationTargetException e) {
Throwable throwable = e.getCause();
while (throwable != null) {
if (throwable instanceof ProgramExecutedException) {
StreamGraph streamGraph = streamExplainEnvironment.getPlan();
}

throwable = throwable.getCause();
}

throw e.getTargetException();
}
} finally {
streamExplainEnvironment.unsetAsContext();
}

1 kubernetes

架构图

1.1 pod

Pod 是 Kubernetes 定义的“容器运行模型”,Docker(或 containerd)只是把模型落地成 Linux 进程的工具。 Pod 运行时不安装操作系统;容器只携带用户态环境,所有 Pod 都运行在宿主机同一个 Linux 内核之上。

docker runc 使用的是:Linux namespaces 或 cgroups

Pod 可以认为是容器的封装,一个 Pod 中可以存在一个或者多个容器

1.2 Namespace

Namespace 是 kubernetes 系统中的一种非常重要资源,它的主要作用是用来实现多套环境的资源隔离或者多租户的资源隔离

默认情况下,kubernetes 集群中的所有的 Pod 都是可以相互访问的。但是在实际中,可能不想让两个 Pod 之间进行互相的访问,那此时就可以将两个 Pod 划分到不同的 namespace 下。kubernetes 通过将集群内部的资源分配到不同的 Namespace 中,可以形成逻辑上的”组”,以方便不同的组的资源进行隔离使用和管理。

1.3 Deployment

2 Flink on k8s

https://juejin.cn/post/7038404897628225573

https://juejin.cn/post/7199235461045960759

https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-main/docs/concepts/architecture/

2.1.1 Application模式

2.1.2 Session模式

2.1.3 峰峦(pre-job模式)

模式 原理 优点 缺点 对我们
Application模式 每个作业各有一个JM 1. 资源隔离,负载均衡保证
2. 作业间彼此独立,互不影响
Session模式 所有作业共享一个JM 资源隔离差,作业间相互影响
峰峦Pre-job模式

2.2 Kubernetes vs Yarn

what is Kubernetes API? how to use?
what is 峰峦?

2.3 Oceanus/Flink调度峰峦资源

2.4 Kubernetes调度过程

K8S不仅能将用户提供的单个容器运行起来,将其对外暴露出去提供服务。还提供了:路由网关、集群监控、灾难恢复,以及应用的水平扩展等能力。

2.5 [参考文献]

  1. 一文讲明白K8S各核心架构组件
  2. 【2022年度技术突破奖】弹性·峰峦云原生大数据平台项目团队

1 Flink实践之延迟统计

流计算是基于消息触发计算的,若没有消息到达到则无法计算,这类指标恰好是要求在指定的超时时间计算出有多少未达到的消息。

1.1 常见场景

1.1.1 菜鸟-物流单配送超时统计

实时数仓的建设难度: 链路复杂、实操节点多、汇总维度多、考核逻辑复杂的特点

仓配实时数据已覆盖了绝大多数场景,但是有这样一类特殊指标:“晚点超时指标”(例如:出库超 6 小时未揽收的订单量),仍存在实时汇总计算困难。
这类指标对于指导实操有着重要意义,可以告知运营小二当前多少订单积压在哪些作业节点,应该督促哪些实操人员加快作业,这对于物流的时效 KPI 达成至关重要。

1.1.2 财付通-资金流账龄超时监控

1.2 解决方案

1.3 方案一: MQ延迟发送

菜鸟的解决方案图

graph LR

T1[订单表]
T2[清洗表]
T3[订单延迟表]
T4[结果表]

J1(任务一: 数据清洗)
J2(任务二: 延迟6小时)
J3(任务三: 汇总计算)

T1 --> J1
J1 --> T2
T2 --> J2
J2 --> T3
T3 --> J3
T2 --> J3
J3 --> T4

简化版:

graph LR

T1[订单表]
T3[Datagen 定时数据表
每秒生成一条数据:
以当前时间作为事件时间] T4[结果表
需支持数据回撤] J1(清洗/过滤) T1 --> J1 subgraph 计算任务 J2(指定事件时间) J3(T1 left join T3
) J1 --> J2 T3 --> J3 J2 --> J3 end J3 --> T4

1.4 方案二: FlinkState+TimeService

1.6 方案四: OLAP


1.7 [参考文献]:

https://flink-learning.org.cn/activity/detail/0ab621d17a16e502dcc6e5e2980e643c?tab=ce84125f7d83480f7ccf34685cce43f6&city=&page=xiangguanhuodong

https://files.alicdn.com/tpsservice/0f9d6f504576175dbf6d084c3e2a43fa.pdf?spm=a2csy.flink.0.0.7f395badaF3TXC&file=0f9d6f504576175dbf6d084c3e2a43fa.pdf

利用Flink实现实时超时统计场景

  1. 利用 blink+MQ 实现流计算中的超时统计问题

17张图带你彻底搞懂hudi upsert源码

如果要深入了解apache hudi技术的应用或是性能调优,那么明白源码中的原理对我们会有很大的帮助。在apache hudi 中upsert 是他的核心功能之一,主要完成增量数据在hdfs上的修改,并可以支持事务。在hive中修改数据需要重新分区或重新整个表,但是对于hudi更新可以是文件级别的重写或是数据先进行追加后续再重写,对比Hive 大大地提高了更新性能。upsert支持两种模式的写入copy on write和merge on read ,下面本文将介绍Apache Hudi 在spark中upsert的内核原理。

对于hudi upsert 操作整理了比较核心的几个操作如图:

在这里插入图片描述

  1. 构造HoodieRecord Rdd对象:hudi 会根据元数据信息构造HoodieRecord Rdd 对象,方便后续数据去重和数据合并。

  2. 数据去重:一批增量数据中可能会有重复的数据,hudi会根据主键对数据进行去重避免重复数据写入hudi 表。

  3. 数据fileId位置信息获取:在修改记录中可以根据索引获取当前记录所属文件的fileid,在数据合并时需要知道数据update操作向那个fileId文件写入新的快照文件。

  4. 数据合并:hudi 有两种模式cow和mor。在cow模式中会重写索引命中的fileId快照文件;在mor 模式中根据fileId 追加到分区中的log 文件。

  5. 完成提交:在元数据中生成xxxx.commit文件,只有生成commit 元数据文件,查询引擎才能根据元数据查询到刚刚upsert 后的数据。

  6. compaction压缩:主要是mor 模式中才会有,他会将mor模式中的xxx.log 数据合并到xxx.parquet 快照文件中去。

  7. hive元数据同步:hive 的元素数据同步这个步骤需要配置非必需操作,主要是对于hive 和presto 等查询引擎,需要依赖hive 元数据才能进行查询。所以在hive 中的同步就是构造外表提供查询。

在这里插入图片描述

介绍完hudi的upsert运行流程,在来看下hudi如何进行存储并且保证事务,在每次upsert完成后都会产生commit 文件记录每次重新的快照文件。

例如上图时间1初始化写入三个分区文件分别是xxx-1.parquet,在时间3场景如果会修改分区1和分区2xxx-1.parquet的数据,那么写入完成后会生成新的快照文件分别是分区1和分区2xxx-3.parquet文件。(上述是cow模式的过程,而对于MOR模式的更新会生成log文件,如果log文件存在追加数据)。如果时间5在去修改分区1的数据那么同理会生成分区1下的新快照文件。可以看出对于hudi 每次修改都是会在文件级别重新写入数据快照。查询的时候就会根据最后一次快照元数据加载每个分区小于等于当前的元数据的parquet文件。hudi事务的原理就是通过元数据mvcc多版本控制写入新的快照文件,在每个时间阶段根据最近的元数据查找快照文件。因为是重写数据所以同一时间只能保证一个事务去重写parquet 文件。不过当前hudi版本加入了并发写机制,原理是zookeeper分布锁控制或者HMS提供锁的方式, 会保证同一个文件的修改只有一个事务会写入成功。

下面将根据spark 调用write方法开始剖析upsert操作每个步骤的执行流程。

2.1 开始提交&数据回滚

在构造好spark 的rdd 后会调用 df.write.format("hudi") 方法执行数据的写入,实际会调用HudiHoodieSparkSqlWriter#write方法实现。在执行任务前hudi 会创建HoodieWriteClient 对象,并构造HoodieTableMetaClient调用startCommitWithTime方法开始一次事务。在开始提交前会获取hoodie 目录下的元数据信息,判断上一次写入操作是否成功,判断的标准是上次任务的快照元数据有xxx.commit后缀的元数据文件。如果不存在那么hudi 会触发回滚机制,回滚是将不完整的事务元数据文件删除,并新建xxx.rollback元数据文件。如果有数据写入到快照parquet 文件中也会一起删除。

在这里插入图片描述

2.2 构造HoodieRecord Rdd 对象

HoodieRecord Rdd 对象得构造先是通过map 算子,提取spark dataframe中的schema和数据,构造avro的genericRecords Rdd, 然后hudi会在进行map算子封装为HoodierRecord Rdd。对于HoodileRecord Rdd 主要由currentlocation,newlocation,hoodiekey,data 组成。HoodileRecord数据结构是为后续数据去重和数据合并时提供基础。

在这里插入图片描述

  • currentLocation 当前数据位置信息:只有数据在当前hudi表中存在才会有,主要存放parquet文件的fileId,构造时默认为空,在查找索引位置信息时被赋予数据。
  • newLocation 数据新位置信息:与currentLocation不同不管是否存在都会被赋值,newLocation是存放当前数据需要被写入到那个fileID文件中的位置信息,构造时默认为空,在merge阶段会被赋予位置信息。
  • HoodieKey 主键信息:主要包含recordKey 和patitionPath 。recordkey 是由hoodie.datasource.write.recordkey.field 配置项根据列名从记录中获取的主键值。patitionPath 是分区路径。hudi 会根据hoodie.datasource.write.partitionpath.field 配置项的列名从记录中获取的值作为分区路径。
  • data 数据:data是一个泛型对象,泛型对象需要实现HoodieRecordPayload类,主要是实现合并方法和比较方法。默认实现OverwriteWithLatestAvroPayload类,需要配置hoodie.datasource.write.precombine.field配置项获取记录中列的值用于比较数据大小,去重和合并都是需要保留值最大的数据。

2.3 数据去重

在upsert 场景中数据去重是默认要做的操作,如果不进行去重会导致数据重复写入parquet文件中。当然upsert 数据中如果没有重复数据是可以关闭去重操作。配置是否去重参数为hoodie.combine.before.upsert,默认为true开启。

在这里插入图片描述

在spark client调用upsert 操作是hudi会创建HoodieTable对象,并且调用upsert 方法。对于HooideTable 的实现分别有cor和mor 两种模式的实现。但是在数据去重阶段和索引查找阶段的操作都是一样的。调用HoodieTable upsert方法后底层的实现都是spark AbstractWriteHelper。在去重操作中,会先使用map 算子提取HoodieRecord中的HoodieatestAvroPayload的实现是保留时间戳最大的记录。**这里要注意如果我们配置的是全局类型的索引map 中的key 值是 HoodieKey 对象中的recordKey。**因为全局索引是需要保证所有分区中的主键都是唯一的,避免不同分区数据重复。当然如果是非分区表,没有必要使用全局索引。

2.4 数据位置信息索引查找

对于hudi 索引主要分为两大类:

  • **非全局索引:**索引在查找数据位置信息时,只会检索当前分区的索引,索引只保证当前分区内数据做upsert。如果记录的分区值发生变更就会导致数据重复。
  • **全局索引:**顾名思义在查找索引时会加载所有分区的索引,用于定位数据位置信息,即使发生分区值变更也能定位数据位置信息。这种方式因为要加载所有分区文件的索引,对查找性能会有影响(hbase 索引除外)。

在这里插入图片描述

spark 索引实现主要有如下几种:布隆索引(BloomIndex),全局布隆索引(GlobalBloomIndex),简易索引(SimpleIndex),简易全局索引(GlobalSimpleIndex),简易全局索引(GlobalSimpleIndex),hbase 索引(HbaseIndex),内存索引(InMemoryHashIndex)。下面会对索引的实现方式做一一介绍。

2.4.1 布隆索引(BloomIndex)

spark 布隆索引的实现类是SparkHoodieBloomIndex ,要想知道布隆索引需要了解下布隆算法一般用于判断数据是否存在的场景,在Hudi中用来判断数据在parquet 文件中是否存在。其原理是计算RecordKey的hash值然后将其存储到bitmap中去,key值做hash可能出现hash 碰撞的问题,为了较少hash 值的碰撞使用多个hash算法进行计算后将hash值存入BitMap,一般三次hash最佳,算法详细参考《漫画:什么是布隆算法?》。但是索引任然有hash 碰撞的问题,但是误判率极低可以保证99%以上的数据判断是正确的。即使有个别数据发生了误判也没有关系在合并操作中如果匹配不到也会被丢弃。

在这里插入图片描述

索引实现类调用tagLocation开始查找索引记录存在哪个parquet 文件中,步骤如下

1.提取所有的分区路径和主键值,然后计算每个分区路径中需要根据主键查找的索引的数量。

2.有了需要加载的分区后,调用LoadInvolvedFiles 方法加载分区下所有的parquet 文件。在加载paquet文件只是加载文件中的页脚信息,页脚存放的有布隆过滤器、记录最小值、记录最大值。对于布隆过滤器其实是存放的是bitmap序列化的对象。

3.加载好parquet 的页脚信息后会根据最大值和最小值构造线段树。

4.根据Rdd 中RecordKey 进行数据匹配查找数据属于那个parqeut 文件中,对于RecordKey查找只有符合最大值和最小值范围才会去查找布隆过滤器中的bitmap ,RecordKey小于最小值找左子树,RecordKey大于最大值的key找右子树。递归查询后如果查找到节点为空说明RecordKey在当前分区中不存在,当前Recordkey是新增数据。查找索引时spark会自定义分区避免大量数据在一个分区查找导致分区数据倾斜。查找到RecordKey位置信息后会构造<HoodieKey,HoodieRecordLocation>Rdd 对象。

在这里插入图片描述

5.以 Rdd 为左表和<HoodieKey,HoodieRecordLocation>Rdd 做左关联,提取HoodieRecordLocation位置信息赋值到HoodieRecord 的currentLocation变量上,最后得到新的HoodieRecord Rdd。在HoodieRecordLocation对象中包含文件fileID 和快照时间。

说明下parquet 文件名称组成,主要是36位fileId、文件编号、spark 任务自定义分区编号、spark 任务stage 编号、spark 任务attempt id、commit时间。

2.4.2 全局布隆索引(GlobalBloomIndex)

对于全局布隆索引底层算法和布隆索引是一样的,只是全局布隆索引在查找每个RecordKey 属于那个parquet 文件中,会加载所有parquet文件的页脚信息构造线段树,然后在去查询索引。因为hudi需要加载所有的parquet文件和线段树节点变多对于查找性能会比普通的布隆索引要差。但是对于分区字段的值发生了修改,如果还是使用普通的布隆索引会导致在当前分区查询不到当成新增数据写入hudi表。这样我们的数据就重复了,在很多业务场景是不被允许的。所以在选择那个字段做分区列时,尽量选择列值永远不会发生变更的,这样我们使用普通布隆索引就可以了。

在这里插入图片描述

全局布隆的实现是继承布隆索引的实现,重写了索引数据的加载和HoodieRecord Rdd左关联部分。加载所有文件的页脚信息刚刚已经提到过了,数据不在知道在哪个分区所以要加载全部文件进行判断。在左关联操作中与普通布隆索引不同的是,如果分区发生了变更,默认情况下会修改HoodieKey 中的partitionPath,数据是不会写到变更后的分区路径下,而是会重写到之前的分区路径下,但是数据的内容还是会更新。如果希望删除旧分区数据,新数据写入当前记录的分区。需要设置hoodie.bloom.index.update.partition.path配置项为true,允许分区数据变更到其他分区。此时,关联查询会在原有的基础上在生成一条删除记录,删除记录HoodieRecordPayload的实现类是EmptyHoodieRecordPayload。EmptyHoodieRecordPayload只会存放hoodieKey的主键信息,在数据合并时会被忽略,达到数据硬删除的目的。这里可以根据业务场景选择是否开启分区变更。

2.4.3 简易索引(SimpleIndex)

简易索引与布隆索引的不同是直接加载分区中所有的parquet数据然后在与当前的数据比较是否存在。这个实现比较简单,我想也是因为这样才这么命名的。以下是简易索引的执行步骤:

1.提取所有的分区路径和主键值。

2.根据分区路径加载所有涉及分区路径的parquet文件的数据主要是HooieKey和fileID两列的数据,构造<HoodieKey,HoodieRecordLocation> Rdd 对象。

3.同布隆索引一样以 Rdd 为左表和<HoodieKey,HoodieRecordLocation>Rdd 做左关联,提取HoodieRecordLocation位置信息赋值到HoodieRecord 的currentLocation变量上,最后得到新的HoodieRecord Rdd

在这里插入图片描述

2.4.4 简易全局索引(GlobalSimpleIndex)

简易全局索引同布隆全局索引一样,需要加载所有分区的parquet 文件数据,构造<HoodieKey,HoodieRecordLocation>Rdd然后后进行关联。在简易索引中hoodie.simple.index.update.partition.path配置项也是可以选择是否允许分区数据变更。数据文件比较多数据量很大,这个过程会很耗时。

2.4.5 hbase 索引(HbaseIndex)

​ hbase 索引与布隆索引、简易索引不同,本身就是一个全局索引,那种场景都可以用。但是需要额外hbase 服务来存储hudi的索引信息,一旦hbase 出现故障hudi upsert 是会无法工作的。不管在布隆索引或简易索引中索引是和parquet 文件是一体的,要么一起成功,要么一起失败。但是在hbase索引中文件和索引是分开在特定的情况下可能有一致性的问题(什么特定情况?)。不过hbase索引最大程度的避免了这一问题。

在这里插入图片描述

hbase 索引实现步骤如下:

1.连接hbase 数据库。

2.批量请求hbase 数据库。

3.检查get 获取的数据是否为有效索引,这时hudi 会连接元数据检查commit时间是否有效,如果无效currentLocation将不会被赋值。检查是否为有效索引的目的是当索引更新一半hbase 宕机导致任务失败,保证不会加载过期索引。避免hbase 索引和数据不一致导致数据进入错误的分区。

4.检查是否开启允许分区变更,这里的做法和全局布隆索引、全局简易索引的实现方式一样。

在hudi中使用hbase索引需要提前建表,hbase表的列簇为_s。示例如下

1
create 'table_name','_s'
1
2
3
4
5
6
7
8
hoodie.index.hbase.zkquorum   必填项:zk连接地址
hoodie.index.hbase.zkport 必填项:zk连接端口
hoodie.index.hbase.zknode.path 必填项:zookeeper znode路径
hoodie.index.hbase.table 必填项:hbase中的表名
hoodie.index.hbase.get.batch.size 默认值100 :每批量请求大小
hoodie.hbase.index.update.partition.path 默认值false: 是否允许分区变更
hoodie.index.hbase.put.batch.size.autocompute 默认值false :是否开启自动计算每个批次大小
hoodie.index.hbase.rollback.sync 默认值false:rollback阶段是否开启同步索引。如果设置为true在写入hbase索引导致hbase 宕机或者jvm oom任务失败,在触发rollback 阶段 会删除失败任务的索引保证索引和数据一致。在上次任务失败且数据分区字段值反复变更时可以避免数据重复。

2.4.6 内存索引(InMemoryHashIndex)

内存索引目前spark 的实现只是构造的一个ConcurrentMap在内存中,不会加载parquet 文件中的索引,当调用tagLocation方法会在map 中判断key值是否存在。spark 内存索引当前只是用来测试的索引。

2.4.7 索引的选择

普通索引:主要用于非分区表和分区不会发生分区列值变更的表。当然如果你不关心多分区主键重复的情况也是可以使用。他的优势是只会加载upsert数据中的分区下的每个文件中的索引,相对于全局索引需要扫描的文件少。并且索引只会命中当前分区的fileid 文件,需要重写的快照也少相对全局索引高效。但是某些情况下我们的设置的分区列的值就是会变那么必须要使用全局索引保证数据不重复,这样upsert 写入速度就会慢一些。其实对于非分区表他就是个分区列值不会变且只有一个分区的表,很适合普通索引,如果非分区表硬要用全局索引其实和普通索引性能和效果是一样的。

全局索引:分区表场景要考虑分区值变更,需要加载所有分区文件的索引比普通索引慢。

布隆索引:加载fileid 文件页脚布隆过滤器,加载少量数据数据就能判断数据是否在文件存在。缺点是有一定的误判,但是merge机制可以避免重复数据写入。parquet文件多会影响索引加载速度。适合没有分区变更和非分区表。主键如果是类似自增的主键布隆索引可以提供更高的性能,因为布隆索引记录的有最大key和最小key加速索引查找。

全局布隆索引:解决分区变更场景,原理和布隆索引一样,在分区表中比普通布隆索引慢。

简易索引:直接加载文件里的数据不会像布隆索引一样误判,但是加载的数据要比布隆索引要多,left join 关联的条数也要比布隆索引多。大多数场景没布隆索引高效,但是极端情况命中所有的parquet文件,那么此时还不如使用简易索引加载所有文件里的数据进行判断。

全局简易索引:解决分区变更场景,原理和简易索引一样,在分区表中比普通简易索引慢。建议优先使用全局布隆索引。

hbase索引:不受分区变跟场景的影响,操作算子要比布隆索引少,在大量的分区和文件的场景中比布隆全局索引高效。因为每条数据都要查询hbase ,upsert数据量很大会对hbase有负载的压力需要考虑hbase集群承受压力,适合微批分区表的写入场景 。在非分区表中数量不大文件也少,速度和布隆索引差不多,这种情况建议用布隆索引。

内存索引:用于测试不适合生产环境

2.5 数据合并

​ cow模式 和mor模式在前面的操作都是一样的,不过在合并的时候hudi构造的执行器就不一样了。对于cow 会根据位置信息中fileId 重写parquet文件,在重写中如果数据是更新会比较parquet文件的数据和当前的数据的大小进行更新,完成更新数据和插入数据。而mor模式会根据fileId 生成一个log 文件,将数据直接写入到log文件中,如果fileID的log文件已经存在,追加数据写入到log 文件中。与cow 模式相比少了数据比较的工作所以性能要好,但是在log 文件中可能保存多次写有重复数据在读log数据时候就不如cow模式了。还有在mor模式中log文件和parquet 文件都是存在的,log 文件的数据会达到一定条件和parqeut 文件合并。所以mor有两个视图,ro后缀的视图是读优化视图(read-optimized)只查询parquet 文件的数据。rt后缀的视图是实时视图(real-time)查询parquet 和log 日志中的内容。

2.5.1 copy on write模式

cow模式数据合并实现逻辑调用BaseSparkCommitActionExecutor类的excute方法,实现步骤如下:

在这里插入图片描述

1.通过countByKey 算子提取分区路径和文件位置信息并统计条数,用于后续根据分区文件写入的数据量大小评估如何分桶。

2.统计完成后会将结果写入到workLoadProfile 对象的map 中,这个时候已经完成合并数据的前置条件。hudi会调用saveWorkloadProfileMetadataToInfilght 方法写入infight标识文件到.hoodie元数据目录中。在workLoadProfile的统计信息中套用的是类似双层map数据结构, 统计是到fileid 文件级别。

3.根据workLoadProfile统计信息生成自定义分区 ,这个步骤就是分桶的过程。首先会对更新的数据做分桶,因为是更新数据在合并时要么覆盖老数据要么丢弃,所以不存在parquet文件过于膨胀,这个过程会给将要发生修改的fileId都会添加一个桶。然后会对新增数据分配桶,新增数据分桶先获取分区路径下所有的fileid 文件, 判断数据是否小于100兆。小于100兆会被认为是小文件后续新增数据会被分配到这个文件桶中,大于100兆的文件将不会被分配桶。获取到小文件后会计算每个fileid 文件还能存少数据然后分配一个桶。如果小文件fileId 的桶都分配完了还不够会根据数据量大小分配n个新增桶。最后分好桶后会将桶信息存入一个map 集合中,当调用自定义实现getpartition方法时直接向map 中获取。所以spark在运行的时候一个桶会对应一个分区的合并计算。

1
2
3
4
5
6
hoodie.parquet.small.file.limit   默认104857600(100兆):小于100兆的文件会被认为小文件,有新增数据时会被分配数据插入。
hoodie.copyonwrite.record.size.estimate 默认1024 (1kb): 预估一条数据大小多大,用来计算一个桶可以放多少条数据。
hoodie.record.size.estimation.threshold 默认为1: 数据最开始的时候parquet文件没有数据会去用默认的1kb预估一条数据的大小,如果有fileid的文件大小大于 (hoodie.record.size.estimation.threshold*hoodie.parquet.small.file.limit) 一条记录的大小将会根据(fileid文件大小/文件的总条数)来计算,所以这里是一个权重值。
hoodie.parquet.max.file.size 默认120 * 1024 * 1024(120兆):文件的最大大小,在分桶时会根据这个大小减去当前fileId文件大小除以预估每条数据大小来计算当前文件还能插入多少数据。因为每条数据大小是预估计算平均值的,所以这里最大文件的大小控制只能接近与你所配置的大小。
hoodie.copyonwrite.insert.split.size 默认500000 :精确控制一个fileid文件存放多少条数据,前提必须关闭hoodie.copyonwrite.insert.auto.split 自动分桶。
hoodie.copyonwrite.insert.auto.split 默认true : 是否开启自动分桶。

4.分桶结束后调用handleUpsertPartition合并数据。首先会获取map 集合中的桶信息,桶类型有两种新增和修改两种。如果桶fileid文件只有新增数据操作,直接追加文件或新建parquet文件写入就好,这里会调用handleInsert方法。如果桶fileid文件既有新增又有修改或只有修改一定会走handUpdate方法。这里设计的非常的巧妙对于新增多修改改少的场景大部分的数据直接可以走新增的逻辑可以很好的提升性能。对于handUpdate方法的处理会先构造HoodieMergeHandle对象初始化一个map集合,这个map集合很特殊拥有存在内存的map集合和存在磁盘的map 集合,这个map集合是用来存放所有需要update数据的集合用来遍历fileid旧文件时查询文件是否存在要不要覆盖旧数据。这里使用内存加磁盘为了避免update桶中数据特别大情况可以将一部分存磁盘避免jvm oom。update 数据会优先存放到内存map如果内存map不够才会存在磁盘map,而内存Map默认大小是1g 。DiskBasedMap 存储是key信息存放的还是record key ,value 信息存放value 的值存放到那个文件上,偏移量是多少、存放大小和时间。这样如果命中了磁盘上map就可以根据value存放的信息去获取hoodieRecord了。

1
2
hoodie.memory.spillable.map.path   默认值 /tmp/ : 存放DiskBasedMap的路径
hoodie.memory.merge.max.size 默认值 1024*1024*1024(1g):内存map的最大容量

5.构造sparkMergHelper 开始合并数据写入到新的快照文件。在SparkMergHelper 内部会构造一个BoundedInMemoryExecutor 的队列,在这个队列中会构造多个生产者和一个消费者(file 文件一般情况只有一个文件所以生产者也会是一个)。producers 会加载老数据的fileId文件里的数据构造一个迭代器,执行的时候先调用producers 完成初始化后调用consumer。而consumer被调用后会比较数据是否存在ExternalSpillableMap 中如果不存在重新写入数据到新的快照文件,如果存在调用当前的HoodileRecordPayload 实现类combineAndGetUpdateValue 方法进行比较来确定是写入老数据还是新数据,默认比较那个数据时间大。这里有个特别的场景就是硬删除,对于硬删除里面的数据是空的,比较后会直接忽略写入达到数据删除的目的。

2.5.2 merge on read模式

在Mor模式中的实现和前面cow模式分桶阶段都是共用的逻辑,这里主要说下最后的合并和cow 模式不一样的操作。在mor 合并是调用AbstarctSparkDeltaCommitActionExecutor 的execute方法,会构造HoodieAppaendHandle 对象。在写入时调用append 向log日志文件追加数据,如果日志文件不存在将新建log文件。

1
2
hoodie.logfile.max.size   默认值:1024 * 1024 * 1024(1g) 日志文件最大大小
hoodie.logfile.data.block.max.size 默认值:256 * 1024 * 1024(256兆) 写入多少数据后刷一次磁盘

分桶相关参数与cow模式通用

在这里插入图片描述

2.6 索引更新

在这里插入图片描述

​ 数据写入到log文件或者是parquet 文件,这个时候需要更新索引。简易索引和布隆索引对于他们来说索引在parquet文件中是不需要去更新索引的。这里索引更新只有hbase索引 和内存索引需要更新。内存索引是更新通过map 算子写入到内存map上,hbase索引通过map算子put到hbase上。

2.7 完成提交

2.7.1 提交&元数据信息归档

上述操作如果都成功且写入时writeStatus中没有任何错误记录,hudi 会进行完成事务的提交和元数据归档操作,步骤如下:

1.sparkRddWriteClient 调用commit 方法,首先会向Hdfs 上提交一个.commit 后缀的文件,里面记录的是writeStatus的信息包括写入多少条数据、fileID快照的信息、Schema结构等等。当commit 文件写入成功就意味着一次upsert 已经成功,hudi 内的数据就可以查询到。

2.为了不让元数据一直增长下去需要对元数据做归档操作。元数据归档会先创建HoodieTimelineArchiveLog对象,通过HoodieTableMetaClient 获取.hoodie目录下所有的元数据信息,根据这些元数据信息来判断是否达到归档条件。如果达到条件构造HooieLogFormatWrite对象对archived文件进行追加。每个元数据文件会封装成 HoodieLogBlock 对象批量写入。

1
2
3
hoodie.keep.max.commits   默认30:最多保留多少个commit元数据,默认会在第31个commit的时候会触发一次元数据归档操作,由这个参数来控制元数据归档时机。
hoodie.keep.min.commits 默认20: 最少保留多少个commit元数据,默认会将当前所有的commimt提交的个数减去20,剩下的11个元数据被归档,这个参数间接控制每次回收元数据个数。
hoodie.commits.archival.batch 默认10 :每多少个元数据写入一次到archived文件里,这里就是一个刷盘的间隔。

在这里插入图片描述

2.7.2 数据清理

元数据清理后parquet 文件也是要去清理,在hudi 有专有spark任务去清理文件。因为是通过spark 任务去清理文件也有对应XXX.clean.request、xxx.clean.infight、xxx.clean元数据来标识任务的每个任务阶段。数据清理步骤如下:

1.构造baseCleanPanActionExecutor 执行器,并调用requestClean方法获取元数据生成清理计划对象HoodieCleanPlan。判断HoodieCleanPlan对象满足触发条件会向元数据写入xxx.request 标识,表示可以开始清理计划。

2.生成执行计划后调用baseCleanPanActionExecutor 的继承类clean方法完成执行spark任务前的准备操作,然后向hdfs 写入xxx.clean.inflight对象准备执行spark任务。

3.spark 任务获取HoodieCleanPlan中所有分区序列化成为Rdd并调用flatMap迭代每个分区的文件。然后在mapPartitions算子中调用deleteFilesFunc方法删除每一个分区的过期的文件。最后reduceBykey汇总删除文件的结果构造成HoodieCleanStat对象,将结果元数据写入xxx.clean中完成数据清理。

1
2
3
4
5
hoodie.clean.automatic  默认true :是否开启自动数据清理,如果关闭upsert 不会执行清理任务。
hoodie.clean.async 默认false: 是否异步清理文件。开启异步清理文件的原理是开启一个后台线程,在client执行upsert时就会被调用。
hoodie.cleaner.policy 默认 HoodieCleaningPolicy.KEEP_LATEST_COMMITS :数据清理策略参数,清理策略参数有两个配置KEEP_LATEST_FILE_VERSIONS和KEEP_LATEST_COMMITS。
hoodie.cleaner.commits.retained 默认10 :在KEEP_LATEST_COMMITS策略中配置生效,根据commit提交次数计算保留多少个fileID版本文件。因为是根据commit提交次数来计算,参数不能大于hoodie.keep.min.commits(最少保留多少次commmit元数据)。
hoodie.cleaner.fileversions.retained 默认3 :在KEEP_LATEST_FILE_VERSIONS策略中配置生效,根据文件版本数计算保留多少个fileId版本文件。

在这里插入图片描述

2.7.3 数据压缩

数据压缩是mor 模式才会有的操作,目的是让log文件合并到新的fileId快照文件中。因为数据压缩也是spark 任务完成的,所以在运行时也对应的xxx.compaction.requet、xxx.compaction.clean、xxx.compaction元数据生成记录每个阶段。数据压缩实现步骤如下:

1.sparkRDDwirteClient 调用compaction方法构造BaseScheduleCompationActionExecutor对象并调用scheduleCompaction方法,计算是否满足数据压缩条件生成HoodieCompactionPlan执行计划元数据。如果满足条件会向hdfs 写入xxx.compation.request元数据标识请求提交spark任务。

2.BaseScheduleCompationActionExecutor会调用继承类SparkRunCompactionExecutor类并调用compact方法构造HoodieSparkMergeOnReadTableCompactor 对象来实现压缩逻辑,完成一切准备操作后向hdfs写入xxx.compation.inflight标识。

3.spark任务执行parallelize加载HooideCompactionPlan 的执行计划,然后调用compact迭代执行每个分区中log的合并逻辑。在 compact会构造HoodieMergelogRecordScanner 扫描文件对象,加载分区中的log构造迭代器遍历log中的数据写入ExtemalSpillableMap。这个ExtemalSpillableMap和cow 模式中内存加载磁盘的map 是一样的。至于合并逻辑是和cow模式的合并逻辑是一样的,这里不重复阐述都是调用cow模式的handleUpdate方法。

4.完成合并操作会构造writeStatus结果信息,并写入xxx.compaction标识到hdfs中完成合并操作。

1
2
3
4
hoodie.compact.inline  默认false:是否在一个事务完成后内联执行压缩操作,这里开启并不一定每次都会触发索引操作后面还有策略判断。
hoodie.compact.inline.trigger.strategy 默认CompactionTriggerStrategy.NUM_COMMITS: 压缩策略参数。该参数有NUM_COMMITS、TIME_ELAPSED、NUM_AND_TIME、NUM_OR_TIME。NUM_COMMITS根据提交次数来判断是否进行压缩;TIME_ELAPSED根据实际来判断是否进行压缩;NUM_AND_TIME 根据提交次数和时间来判断是否进行压缩;NUM_OR_TIME根据提交次数或时间来判断是否进行压缩。
hoodie.compact.inline.max.delta.commits 默认5 :设置提交多少次后触发压缩策略。在NUM_COMMITS、NUM_AND_TIME和NUM_OR_TIME策略中生效。
hoodie.compact.inline.max.delta.seconds 默认60 * 60(1小时):设置在经过多长时间后触发压缩策略。在TIME_ELAPSED、NUM_AND_TIME和NUM_OR_TIME策略中生效。

在这里插入图片描述

2.8 hive元数据同步

​ 实现原理比较简单就是根据hive外表和hudi表当前表结构和分区做比较,是否有新增字段和新增分区如果有会添加字段和分区到hive外表。如果不同步hudi新写入的分区无法查询。在cow模式 中只会有ro表(读优化视图,而在mor模式中有ro表(读优化视图)和rt表(实时视图)。

2.9 提交成功通知回调

​ 当事务提交成功后向外部系统发送通知消息,通知的方式有两种,一种是发送http服务消息,一种是发送kafka 消息。这个通知过程我们可以把清理操作、压缩操作、hive元数据操作,都交给外部系统异步处理或者做其他扩展。也可以自己实现HoodieWriteCommitCallback的接口,自定义实现。

1
2
3
4
5
6
7
8
9
10
11
12
hoodie.write.commit.callback.on   默认false:是否开启提交成功后向外部系统发送回调指令。
hoodie.write.commit.callback.class 默认org.apache.hudi.callback.impl.HoodieWriteCommitHttpCallback: 配置回调实现类,默认通过Http的方式发送消息到外部系统
http实现类配置参数
hoodie.write.commit.callback.http.url 无默认配置项:外部服务http url地址。
hoodie.write.commit.callback.http.api.key 默认hudi_write_commit_http_callback:外部服务http请求头HUDI-CALLBACK-KEY的值,可用于服务请求验签使用。
hoodie.write.commit.callback.http.timeout.seconds 默认3秒:请求超时时间。
kafka实现类配置参数
hoodie.write.commit.callback.kafka.bootstrap.servers 无默认值:配置kafka broker 服务地址。
hoodie.write.commit.callback.kafka.topic 无默认值:配置kafka topic名称。
hoodie.write.commit.callback.kafka.partition 无默认值:配置发送到那个kafka broker分区。
hoodie.write.commit.callback.kafka.acks 默认值all:配置kafka ack。
hoodie.write.commit.callback.kafka.retries 默认值值3:配置kafka 失败重试次数。

在分析spark upsert源码中还有很多细节是略过,如时间线服务、并发写机制、cluster模式、Hfile格式的写入和merge等等。篇幅有限先解析这么多,希望本文能帮你了解spark upsert的内核原理。谢谢大家阅读本文。

1 基础技术栈

1.1 Java

Java 全栈知识体系

1.2 数据结构

1.2.1 LSM Tree

1.2.2 B-Tree

1.3 底层存储引擎

1.3.1 LevelDB

1.3.2 RocksDB

https://github.com/facebook/rocksdb/wiki/
https://alexstocks.github.io/html/rocksdb.html

1.4 硬件与存储

1.4.1 OSS 对象存储

2 分布式协调

2.1 Zookeeper

3 存储层

3.1 HDFS

分布式文件系统,大数据存储的基础设施

4 消息队列

4.1 Kafka

4.2 Pulsar

详情

5 批处理引擎

5.1 Spark

深入理解 Spark(Spark SQL 优化、数据倾斜处理)和 Hive(调优、存储格式如 Parquet/ORC)。

5.1.1 Spark Core

5.1.2 Spark SQL

5.1.3 性能优化

  • 数据倾斜处理
  • 算子优化
  • 资源调优

6 流处理引擎

需要掌握其状态管理(State)、容错机制(Checkpoint)以及 Exactly-once 语义。

详情

6.1.1 Runtime 架构

6.1.1.1 Memory Management

6.1.3 State & Checkpoint

6.1.3.1 状态管理

6.1.3.2 容错机制

6.1.3.3 Exactly-once 语义

6.1.4 RocksDB State Backend

7 湖仓一体架构

存储与湖仓一体: 掌握 Data Lakehouse 架构(如 Iceberg、Hudi、Delta Lake),解决近实时更新和流批统一的问题。

7.1 Apache Iceberg

7.2 Apache Hudi

7.3 Apache Paimon

7.4 Delta Lake

8 OLAP 分析引擎

为了满足秒级响应的分析需求,架构师需要根据场景(并发量、延迟、数据量)选择合适的引擎:

8.1 ClickHouse

极致的单表查询性能

8.2 Apache Doris

优秀的国产自研引擎,支持高并发、流批一体和极简运维

Doris为什么那么快?

8.3 StarRocks

优秀的国产自研引擎,支持高并发、流批一体和极简运维

Apache Doris和Clickhouse的深度分析

8.4 Presto / Trino

适用于跨数据源的联邦查询

9 云原生大数据

9.1 容器化部署

9.2 Kubernetes 调度

9.3 资源弹性伸缩

9.4 Ray

10 数据治理与架构设计

10.1 数据治理与全生命周期管理

架构师的价值往往体现在”如何管好数据”而非仅仅”存下数据”。

10.1.1 元数据管理

构建数据字典,实现血缘追踪(Data Lineage)

10.1.2 数据质量

建立监控体系,涵盖完整性、准确性和一致性

10.1.3 数据安全

  • 权限控制(Ranger)
  • 数据脱敏
  • 数据加密
  • 隐私计算

10.1.4 数仓建模

  • 离线数仓建模
  • 实时数仓建模
  • 维度建模(Kimball)
  • Data Vault

10.2 架构设计思维

这是专家级与高级开发的本质区别:

10.2.1 选型能力

能够在不同的业务场景(如电商大促、金融风控、物联网监控)下给出最合适的方案

10.2.2 高可用与灾备

设计具备容错能力的多活架构

10.2.3 成本优化

通过冷热数据分级存储、资源弹性伸缩(FinOps)降低云账单成本

11 参考资料

Daily-0609~0615

1 列存

列存具有「读友好,写不友好」的特点。这个特点使列存和 AP 数据库仿佛王八看绿豆一样——对上眼了。因为 AP 数据库正好重视读性能(大查询吞吐量),不重视写性能(能容忍 T+1h/1d 更新)。于是,列存顺理成章地霸占了 AP 数据库的存储底座。

近年来,随着业务的发展,越来越多的业务场景霸道地要求 AP 数据库必须具备实时分析能力。为此,AP 数据库必须把上游 TP 数据库产生的写入和更新实时地导入进来。这对 AP 数据库的列存更新效率制造了新的挑战。

本文将介绍迎接这个挑战诞生的列存高效更新技术。探讨内容包括:

  1. 列存更新的难点。

  2. 一种低效更新方案。

  3. 四类高效更新方案,并介绍这些方案在业界数据库中的实现。涉及数据库包括 Iceberg,Hudi,Kudu,Doris,ADB,Hologres 等。

  4. 比较总结所有方案。

希望本文能对读者了解列存更新技术有所帮助。

2 列存更新的难点

我认为主要面临两个难点:一个是写放大,另一个是无法简单地做 in-place update。
Image

首先这里的写放大特指 IO 次数的放大,具体来说:

  • 传统列存(每列一个文件):写入需要每列一次磁盘 IO。尽管可以通过攒批来均摊开销,但大宽表(成百上千列)场景仍然力不从心。

  • 行列混存(PAX 格式):文件被切分成很多 block,一行的所有列数据按列存格式塞到同一个 block 里面。写入需要在内存中攒满一个列存 block,然后以 block 为单位压缩并刷盘。所有列一次 IO 搞定。尽管这种做法优化了磁盘 IO 次数,但在内存 IO 次数方面还是有放大,因为内存 block 不是行存。

下文所有数据库的列存格式均属于行列混存。

其次是无法简单地做 in-place update。原因有二:

  1. 容易造成比较大的写放大。列存的 block 非常大,因为面向读场景优化,追求高压缩率。哪怕是 in-place update 一个 field,也有可能导致大 block 内数据的 reorganize。重写大 block 会带来比较大的写放大。

  2. 最致命的,AP 数据库的趋势都是 share-storage 架构,基于 HDFS/S3。这类分布式存储本身就不支持 in-place update。所以没得选了,只能做 out-of-place update。

下文所有列存更新技术均属于 out-of-place update。

3 低效更新技术

想了解什么是高效更新技术,必须先了解什么是低效更新技术。—— 黄金架构师(知乎和公众号同名)

最简单的 out-of-place update 方案是 file-level COW(copy-on-write) 。这种方案简单来说就是更新时无论更新一行还是一批,都直接把原文件拷贝出来更新,然后生成一个新文件。

File-level COW 方案在业界的典型代表是 Hudi COW 表。
Image

Hudi COW 表的实现概述如下:

  1. 宏观上看,一张表的数据存储在 HDFS/S3 上的很多列存文件中。

  2. 更新时重写整个列存文件。因为要做写写冲突检测(ww-conflict check;为实现 snapshot isolation),更新需要保证可串行化。Hudi COW 通过 file-level 乐观锁来保证这一点。更新文件期间不加锁,commit 时检查有没有并发写事务抢先更新了相同的文件,如果有,那么自己就 abort。换句话说,不支持对同一文件的并发更新。

  3. 读取时,按快照选取适当版本的文件。

这个方案的优点是读性能非常棒,文件级别多版本使得读完全不受写的影响。缺点是写性能很差,因为写放大很大,并且写并发度低。

显然,file-level COW 是一种低效更新方案。如果用这个方案来应对实时场景,那无异于以卵击石。

到现在为止,我们已经具备了低效更新技术的认知。那么,接下来我们就可以往高效更新技术的殿堂进发了。

4 高效更新技术

我们可以思考一下,要想高效,应该做好哪些事情?我认为,主要是两件事情:

  1. 降低写放大,提升写并发。单线程性能靠降低写放大来优化,多线程性能靠提升写并发来优化。单线程和多线程都在手,性能我有。

  2. 尽量少损害读性能。毕竟咱是 AP 数据库,读性能至关重要,还要靠它来吃饭。

为了做好这两件事情,我们应该追求:

  1. 更细粒度的更新。Hudi COW 拷贝更新整个文件算是 file-level update。如果我们能做到 block-level,tuple-level,甚至是 field-level,写放大会显著优化。

  2. 更细粒度的并发。Hudi COW 加 file-level 的锁,算是 file-level concurrency。如果我们能做到 block-level,tuple-level,写并发会显著提升。(没有 filed-level,复杂且开销不一定小)。当然我们也要清醒地认识到,细粒度的并发需要细粒度的锁,细粒度的锁对于批量更新没有那么友好(想象一下更新整张表的数据,为每一行加一个行锁)。因此这里存在一些取舍的空间。

上述思考很有意思。我运用这个思考框架,再结合业界数据库的实际情况,把业界数据库的高效更新方案按照写并发和更新粒度分成了四类,分别是:

  1. Table-level concurrency + tuple-level update.

  2. File-level concurrency + tuple-level update.

  3. Tuple-level concurrency + field-level update.

  4. Unlimited concurrency + tuple-level update. (Unlimited concurrency 看起来非常唬人,其含义我们下文再细说。)

这四类方案在各家数据库中都是怎么实现的呢?接下来我们挨个分析一下具体的案例。

4.1 Table-level concurrency + tuple-level update

这类方案的典型代表是 Iceberg MOR 表。

使用这类方案的数据库绝对地重视批量更新,轻视并发更新和单行更新,因此更新通过表锁来实现,仅支持 table-level concurrency。Image
Iceberg MOR 表的实现概述如下:

  1. 宏观看是两个存储结构:存储主要数据的列存文件(data file)和存储标记删除的文件(delete file)。data file 和它对应的 delete file 会定期合并成新的 data file。

  2. 更新流程:update 看做 delete + insert。先向 delete file 追加一条 delete mark 删除老版本,然后把新版本写入 data file。

  • Iceberg 的 delete mark 有两种:1. position delete:记录删除行在 data file 中的行号,行号需要在更新的时候从 data file 中查到。2. equation delete:记录删除的等值条件。equation delete 的写入更高效,尤其是在根据非主键(唯一键)的等值删除场景。这个场景下一个等值条件会匹配很多行。equation delete 只记录一个等值条件,而 position delete 需要为每行查行号并为每行记录 delete mark。

  • Iceberg 采用 table-level 乐观锁。更新操作完成后,需要把新增的 delete file 文件登记到表的元数据里面去。如果此时发现别的并发事务往这张表已经登记过元数据,那么自己就 abort(某些情况下会挣扎重试一下)。

  • delete mark 和新版本 tuple 加起来的数据量级约等于一个 tuple,所以是 tuple-level update,写放大很小。

  1. 读取流程:需要 merge 一下 data file 和 delete file(这也是 MOR(merge on read) 名称的由来)。
  • 如果是 position delete,只需要按行号顺序地归并 一下(因为行号天然有序)。

  • 如果是 equation delete,需要为每一行计算可能很多(每删除一次就多一个)的等值条件是否匹配。

  • 二者对比,显然 position delete 的读取更高效。

相比 Hudi COW,Iceberg 的并发能力差了些。抛开并发不谈,Iceberg 算是牺牲了一些读取的性能(需要 merge on read),换取更新的性能。在这个基础上,Iceberg 还提供了 position delete 和 equation delete 两种方式,给用户提供了 MOR 模式下进一步在读友好和写友好之间权衡的空间,这个做法很「用户友好」。

4.2 File-level concurrency + tuple-level update

这类方案的典型代表是 Hudi MOR 表。

这类数据库的特点是重视批量更新,但没有 Iceberg 那么极端,所以并发度稍高一点,做到了 file-level。Image
Hudi MOR 表的实现概述如下:

  1. 宏观看是两个存储结构:存储主要数据的列存文件(base file)和存储增量更新的行存文件(log file)。log file 采用行存对写比较友好。base file 和它对应的 log file 会定期合并成新的 base file。

  2. 更新流程:update 看做 delete + insert。先向 log file 插入一条 delete mark 来标记删除老版本,然后再向 log file 写入新版本。

  • delete mark 记的是主键。但 log file 中的数据是 append 的,并不按照主键排序。

  • 更新和 Hudi COW 一样,是 file-level OCC,因此是 file-level concurrency。

  • 写入数据量是 delete mark 加新版本的 tuple,因此是 tuple-level update。

  1. 读取流程:需要 merge 一下 base file 和 log file,实现上是 base file 和 log file 按主键做 hash join。由于 log file 中的数据是无序的,即便通过 zone-map 过滤后只需要读一个 base file 的 block,也要 join 整个 log file。

除了并发度比 Iceberg 高以外,Hudi MOR 和 Iceberg 主要有两个区别。

第一个区别是 delete mark 的实现不同

写性能方面,Iceberg position delete < Hudi 主键 delete < Iceberg equation delete。

position delete < 主键 delete,是因为 AP 数据库的更新大多数情况下是提供了整行数据的 upsert,这种情况下,主键 delete 可以做到不需要去 base file 中读数据,直接写 log file 就搞定,而 position delete 还得去读行号。

主键 delete < equation delete,这个显而易见,equation delete 在根据非主键(唯一键)的等值删除场景具有绝对的优势。

读性能方面,Iceberg position delete > Hudi 主键 delete > Iceberg equation delete。

因为主键 delete 需要 base file 和 log file 做 hash join,得构建 hash table 和按行 probe。而 position delete 只需要按天然有序的行号归并,因此更快。equation delete 需要每一行计算很多的等值条件,因此更慢。

第二个区别是同一主键的新老版本存储位置不同。Hudi 能保证新老版本逻辑上在同一个 base file 中(在 base file 或者它对应的 log file 中),而 Iceberg 新版本可能出现在不同的 data file 中。相比之下,Iceberg 在实现主键行级别的索引,主键 zone-map 过滤方面出于劣势。

综合来看,没有明显的胜者,二者互有胜负。

4.3 Tuple-level concurrency + Field-level update

这类方案的唯一代表是 Kudu。一个数据库自成一类。

这类数据库非常重视单行更新,因此通过行锁来实现并发,实现了 tuple-level concurrency。(注意:我不确定 Kudu 是否还支持表锁,通过 cost 判断该加表锁还是行锁。感兴趣的读者可以自行研究一下。)

这类数据库非常「吝惜」存储空间,因此更新时只记录被更新的 field。

Image
Kudu 的实现概述如下

  1. 宏观看有点像一个只有 L0 层的 LSM-tree,但 Kudu 对于写入和更新分类处理。写入和更新分别使用自己的 memtable 和文件。写入写到行存 MemRowSet(Masstree 实现),定期刷盘成为 DiskRowSet 列存文件。更新写到行存 DeltaMemStore,定期刷盘成为 DiskRowSet 对应的 REDO records 文件。DiskRowSet 文件内嵌一个主键 B-tree 索引,用来加速主键点查。DiskRowSet 和 REDO records 文件会被周期性地合并。

  2. 更新流程:先根据主键查询在 DiskRowSet 中的行号,然后在 DeltaMemStore 中记录行号以及更新后的 field。因为只记了更新后的 field,有点像 redo log,所以文件名叫 REDO records。

  • 通过行锁来保证并发更新的事务性,所以是 tuple-level concurrency。

  • 写入的数据量只是更新后的 filed,所以是 field-level update。

  1. 读取流程:需要按行号顺序地归并 DiskRowSet 和 REDO records。

Kudu 的方案是我最喜欢的方案。这个方案写并发高,写放大最小。此外,field-level update 非常有利于大宽表的部分列更新场景,因为不需要花费大量的 IO(每列一次 IO)去补全其他列数据。如果非要找出一个缺点,那就是行锁对批量更新没那么友好。不过这也可以通过动态地选择表锁或者行锁来优化。

4.4 Unlimited concurrency + tuple-level update

这类方案的典型代表是 Doris 和阿里云 ADB。

我创造的 unlimited concurrency 这个术语有点唬人。其实这个术语翻译成人话就是不支持并发更新的事务性——不加任何锁,不检测写写冲突,允许丢失更新(lost-update)。

比如 Doris,并发更新采用 last writer wins 策略,同一行的并发更新,后来的更新会覆盖前面的更新,造成前面的更新丢失。

至于 ADB,我没有看到 ADB 论文中提到并发更新,找官方文档也没找到(也有点迷路了,ADB 有很多版本,MySQL,Postgres 等等,弄晕我了),所以我猜测它也不支持并发更新的事务性,不然没理由不高调宣传。因此,我将 ADB 也归类到了这里。尽管我不能 100% 确定它属于这个分类,但我认为分类错了也问题不大,它不影响本文的核心思想。

由于不支持更新的事务性,这类数据库的单行更新和批量更新都比较快。

提问:牺牲事务性,换取更新的性能,你认为值得吗?欢迎留言讨论。

Image
Doris MOR(unique-key 表) 的实现方式概述如下:

  1. 宏观上看是一个 LSM-tree,只有 memtable 和 L0 层列存文件。

  2. 更新流程:采用 LSM-tree 的做法,直接写入新版本,覆盖旧版本。

  • 每个版本上带着一个 sequence number,作为版本排序的依据。

  • Doris 默认不支持并发更新(官方文档)。如有需求,可以通过开关打开,但会有丢失更新的异象。因为不支持并发更新的事务性,所以我把 Doris 划分为 unlimited concurrency。

  • 写入数据量只是一个 tuple,所以是 tuple-level update。

  1. 读取流程:需要 merge 多个文件以及 memtable 来找到同一主键的最新版本。

抛开对并发更新事务性的支持不谈,Doris MOR 方案的优点是更新非常快。甚至写入不需要关心是否违反主键唯一约束,因为同一主键的新版本总会覆盖旧版本,保证了绝对的唯一。此外,在提供整行数据的 upsert 场景,Doris 也不需要读老版本数据,直接写新版本就可以搞定。Doris MOR 的缺点是读取的并行度不高,无法做到 file-level 独立并行读。因为文件之间存在依赖,需要多个文件合并到一起,才能找到最新版本的数据。在这个限制下,聚合函数无法下推到文件层并行计算,读性能比较差。

Image
ADB 的实现概述如下:

  1. 宏观看存储是两个结构:磁盘上的列存文件(detail file)和内存中的 delete bitmap。delete bitmap 被切分成很多 compressed segment。

  2. 更新流程:update 看做 delete + insert。先在 delete bitmap 标记删除老版本(把删除行的对应 bit 置 1 表示删除),然后写入新版本。

  • delete bitmap 自身的更新通过 segment-level COW 实现,有一定的内存空间开销。

  • 疑似不支持并发更新的事务性,所以我把它划分到 unlimited concurrency。

  • 写入的数据量是内存 bitmap 的一个 compressed segment,以及磁盘上的一个 tuple。因为我重点关注磁盘的写入量,所以这个方案被分类为 tuple-level update。

  1. 读取需要 merge 列存文件和对应的 delete bitmap segments。

ADB 方案最特别之处在于 in-memory delete bitmap,优点是标记删除的读写性能都很棒,毕竟内存数据结构的威力不可小觑。缺点是内存开销比较大。

相比与 Doris,ADB 可以实现文件级别甚至更细粒度的并行读。虽然同一主键的新老版本也可能会出现在不同的文件中,但每个文件有自己的 delete bitmap,可用来过滤被删除的数据。每个文件都知道自己的数据要么是新版本的,要么是被删除的,不依赖别的文件来确定,所以可直接并行读。

In-memory delete bitmap + segment-level COW 的内存开销很大。一种优化的思路是将其转化为 on-disk bitmap。Hologres 就采用了 on-disk bitmap 的方案,每个文件一个 delete bitmap(roaring bitmap 实现),存储在整张表一个的 LSM-tree 中。为了省空间,delete bitmap 的更新不是通过 COW 来实现,而是通过 MOR 来实现,即每个增量更新都会生成一个新版本的 bitmap,读的时候合并所有版本。感兴趣的可以阅读下 Hologres 的论文。更进一步地,Doris MOW(merge on write) 表在 Hologres 方案的基础上,为了减小合并多个 bitmap 的开销,还会把合并后的结果缓存起来,详情可以参考 Doris MOW 的设计文档。

5 总结

分类 file-level concurrency + file-level update table-level concurrency + tuple-level update file-level concurrency + tuple-level update tuple-level concurrency + field-level update unlimited concurrency + tuple-level update
典型代表 Hudi COW Iceberg Hudi MOR Kudu Doris MOR, ADB
特点 读友好 批量更新友好,并发更新差 批量更新友好,并发更新中 单行更新友好,部分列更新友好,并发更新友好,批量更新一般 单行批量更新都友好,事务支持差(不支持并发更新或并发更新异常)

6 参考

  1. https://www.dremio.com/blog/

  2. https://hudi.apache.org/docs/concepts/

  3. https://kudu.apache.org/kudu.pdf

  4. https://cwiki.apache.org/confluence/display/DORIS

  5. https://www.vldb.org/pvldb/vol12/p2059-zhan.pdf

  6. https://kai-zeng.github.io/papers/hologres.pdf

MySQL JDBC Driver源码解读

You can get the source code from Github: https://github.com/mysql/mysql-connector-j

当调用com.mysql.cj.jdbc.Driver类时, 静态代码块会默认加载驱动.

Driver

DriverManager.getConnection时,会通过SPI获取

无需手动注册

ClassLoader加载com.mysql.cj.jdbc.Driver时,会通过静态内部类注册。
也就是说,SPI扫到这个类,就会注册.

设计比较巧妙的一点:

classDiagram

NonRegisteringDriver <|-- Driver

NonRegisteringDriver: Connection connect(String url, Properties info)

创建Connection

sequenceDiagram
  participant ClientImpl
  participant DriverManager
  participant Driver
  participant Connection
  participant NativeSession
  participant NativeProtocol
  participant StatementImpl

alt 创建连接
ClientImpl ->> DriverManager:  getConnection()
DriverManager ->> DriverManager: 加载MySQL的Driver
DriverManager ->> Driver: connect()
Driver ->> Connection: ()
Connection ->> NativeSession: connect()
NativeSession ->> NativeProtocol: connect()
end

alt 执行查询
ClientImpl ->> Connection:  createStatement()
Connection ->> StatementImpl: ()
Connection -->> ClientImpl:  Statement

ClientImpl ->> StatementImpl:  execute()
StatementImpl ->> NativeSession: execSQL()获取结果
end

alt 查询结果
ClientImpl ->> StatementImpl: getResultSet()
StatementImpl -->> ClientImpl: ResultSet
end

构建ConnectionUrl

数据库Url和连接参数 的容器
根据schema来区分ConnectionUrl的类型ConnectionUrl.Type, Type中定义了创建不同ConnectionUrl的类名和参数。

构建Connection

Connection代表与特定数据库的会话。在 Connection 的上下文中,执行 SQL 语句并返回结果。

没有Connection持有一个NativeSession

创建新的socket链接 -> 与服务端握手 -> 构建与服务端的通信协议 -> 通过MySQL协议与服务端连接

NativeSession

与SocketServer的会话

NativeSocketConnection

持有与服务端链接的InputStream和OutputStream

SocketFactory

允许在驱动程序中创建可插入套接字的接口

创建Socket链接

com.mysql.cj.protocol.Protocol

协议提供与 MySQL 服务器通信的设施

充斥着好多pocket相关的类

DNS更新后新ip未生效问题

StandardSocketFactory.connect方法中通过域名获取ip地址。

https://github.com/mysql/mysql-connector-j/blob/805f872a57875f311cb82487efcfb070411a3fa0/src/main/core-impl/java/com/mysql/cj/protocol/StandardSocketFactory.java#L130

1
InetAddress[] possibleAddresses = InetAddress.getAllByName(this.host);
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// look-up or remove from cache  
Addresses addrs;
if (useCache) {
addrs = cache.get(host);
} else {
addrs = cache.remove(host);
if (addrs != null) {
if (addrs instanceof CachedAddresses) {
// try removing from expirySet too if CachedAddresses
expirySet.remove(addrs);
}
addrs = null;
}
}

if (addrs == null) {
// create a NameServiceAddresses instance which will look up
// the name service and install it within cache... Addresses oldAddrs = cache.putIfAbsent(
host,
addrs = new NameServiceAddresses(host, reqAddr)
);
if (oldAddrs != null) { // lost putIfAbsent race
addrs = oldAddrs;
}
}

// ask Addresses to get an array of InetAddress(es) and clone it
return addrs.get().clone();

#InetAddress 带了缓存

在Java中,InetAddress 类用于处理IP地址,其内部实现了DNS缓存机制以提高域名解析的效率。通过设置JVM参数或者系统属性,可以调整InetAddress的DNS缓存行为。

1. JVM参数设置

可以通过在启动Java应用程序时传递JVM参数来设置DNS缓存的行为。以下是常用的JVM参数:

  • java -Dnetworkaddress.cache.ttl=10 -jar test.jar:这个参数设置了DNS缓存的时间为10秒。默认情况下,如果不解设,缓存时间为30秒。
  • java -Dnetworkaddress.cache.ttl=0 -jar test.jar:设置为0意味着不缓存解析成功的结果,每次解析都会访问DNS服务器。
  • java -Dnetworkaddress.cache.ttl=-1 -jar test.jar:设置为-1表示永久缓存解析成功的结果,直到应用程序重启。

2. 系统属性设置

除了JVM参数,还可以通过设置系统属性来调整DNS缓存策略:

  • Security.setProperty("networkaddress.cache.ttl", "-1"):通过编程方式设置DNS缓存时间为永久。

3. 源码级别的设置

在Java代码中,可以通过修改InetAddressCachePolicy类的行为来控制DNS缓存策略。例如,可以通过修改cachePolicy静态变量的值来控制缓存时间。

4. 安全管理器(SecurityManager)的影响

当设置了安全管理器(SecurityManager)时,对于DNS缓存时间的默认行为会有所改变:

  • 如果设置了SecurityManager,默认的CachePolicyFOREVER,即永久缓存DNS的结果,直到进程终止。
  • 如果没有设置SecurityManager,可以通过修改java.security文件中的相关属性来控制DNS缓存行为。

注意事项

  • 永久缓存DNS结果(FOREVER)可能会导致应用程序无法响应DNS变化,因此在生产环境中应谨慎使用。
  • 缓存策略的设置应该根据应用程序的实际需求和网络环境来决定,以确保既高效又安全。
  • 在设置DNS缓存策略时,还应考虑无效DNS缓存时间的设置,即当DNS解析失败时的缓存时间,默认情况下这个时间为10秒。

总结

  1. 创建对象,大量使用反射的原因?
    灵活