0%

DataLake三剑客

作者:辛庸,阿里巴巴计算平台事业部 EMR 技术专家。Apache Hadoop,Apache Spark contributor。对 Hadoop、Spark、Hive、Druid 等大数据组件有深入研究。目前从事大数据云化相关工作,专注于计算引擎、存储结构、数据库事务等内容。

共同点

定性上讲,三者均为 Data Lake 的数据存储中间层,其数据管理的功能均是基于一系列的 meta 文件。meta 文件的角色类似于数据库的 catalog/wal,起到 schema 管理、事务管理和数据管理的功能。与数据库不同的是,这些 meta 文件是与数据文件一起存放在存储引擎中的,用户可以直接看到。这种做法直接继承了大数据分析中数据对用户可见的传统,但是无形中也增加了数据被不小心破坏的风险。一旦某个用户不小心删了 meta 目录,表就被破坏了,想要恢复难度非常大。

Meta 文件包含有表的 schema 信息。因此系统可以自己掌握 Schema 的变动,提供 Schema 演化的支持。Meta 文件也有 transaction log 的功能(需要文件系统有原子性和一致性的支持)。所有对表的变更都会生成一份新的 meta 文件,于是系统就有了 ACID 和多版本的支持,同时可以提供访问历史的功能。在这些方面,三者是相同的。

下面来谈一下三者的不同。

Hudi

先说 Hudi。Hudi 的设计目标正如其名,Hadoop Upserts Deletes and Incrementals(原为 Hadoop Upserts anD Incrementals),强调了其主要支持 Upserts、Deletes 和 Incremental 数据处理,其主要提供的写入工具是 Spark HudiDataSource API 和自身提供的 DeltaStreamer,均支持三种数据写入方式:UPSERT,INSERT 和 BULK_INSERT。其对 Delete 的支持也是通过写入时指定一定的选项支持的,并不支持纯粹的 delete 接口。

其典型用法是将上游数据通过 Kafka 或者 Sqoop,经由 DeltaStreamer 写入 Hudi。DeltaStreamer 是一个常驻服务,不断地从上游拉取数据,并写入 hudi。写入是分批次的,并且可以设置批次之间的调度间隔。默认间隔为 0,类似于 Spark Streaming 的 As-soon-as-possible 策略。随着数据不断写入,会有小文件产生。对于这些小文件,DeltaStreamer 可以自动地触发小文件合并的任务。

在查询方面,Hudi 支持 Hive、Spark、Presto。

在性能方面,Hudi 设计了 `
HoodieKey
,一个类似于主键的东西。
HoodieKey
有 Min/Max 统计,BloomFilter,用于快速定位 Record 所在的文件。在具体做 Upserts 时,如果 HoodieKey
不存在于 BloomFilter,则执行插入,否则,确认 HoodieKey
是否真正存在,如果真正存在,则执行 update。这种基于 HoodieKey + BloomFilter 的 upserts 方法是比较高效的,否则,需要做全表的 Join 才能实现 upserts。对于查询性能,一般需求是根据查询谓词生成过滤条件下推至 datasource。Hudi 这方面没怎么做工作,其性能完全基于引擎自带的谓词下推和 partition prune 功能。

Hudi 的另一大特色是支持 Copy On Write 和 Merge On Read。前者在写入时做数据的 merge,写入性能略差,但是读性能更高一些。后者读的时候做 merge,读性能查,但是写入数据会比较及时,因而后者可以提供近实时的数据分析能力。

最后,Hudi 提供了一个名为 run_sync_tool 的脚本同步数据的 schema 到 Hive 表。Hudi 还提供了一个命令行工具用于管理 Hudi 表。

hudi
image.png


Iceberg

Iceberg 没有类似的 HoodieKey 设计,其不强调主键。上文已经说到,没有主键,做 update/delete/merge 等操作就要通过 Join 来实现,而 Join 需要有一个 类似 SQL 的执行引擎。Iceberg 并不绑定某个引擎,也没有自己的引擎,所以 Iceberg 并不支持 update/delete/merge。如果用户需要 update 数据,最好的方法就是找出哪些 partition 需要更新,然后通过 overwrite 的方式重写数据。Iceberg 官网提供的 quickstart 以及 Spark 的接口均只是提到了使用 Spark dataframe API 向 Iceberg 写数据的方式,没有提及别的数据摄入方法。至于使用 Spark Streaming 写入,代码中是实现了相应的 StreamWriteSupport,应该是支持流式写入,但是貌似官网并未明确提及这一点。支持流式写入意味着有小文件问题,对于怎么合并小文件,官网也未提及。我怀疑对于流式写入和小文件合并,可能 Iceberg 还没有很好的生产 ready,因而没有提及(纯属个人猜测)。

在查询方面,Iceberg 支持 Spark、Presto。

Iceberg 在查询性能方面做了大量的工作。值得一提的是它的 hidden partition 功能。Hidden partition 意思是说,对于用户输入的数据,用户可以选取其中某些列做适当的变换(Transform)形成一个新的列作为 partition 列。这个 partition 列仅仅为了将数据进行分区,并不直接体现在表的 schema 中。例如,用户有 timestamp 列,那么可以通过 hour(timestamp) 生成一个 timestamp_hour 的新分区列。timestamp_hour 对用户不可见,仅仅用于组织数据。Partition 列有 partition 列的统计,如该 partition 包含的数据范围。当用户查询时,可以根据 partition 的统计信息做 partition prune。

除了 hidden partition,Iceberg 也对普通的 column 列做了信息收集。这些统计信息非常全,包括列的 size,列的 value count,null value count,以及列的最大最小值等等。这些信息都可以用来在查询时过滤数据。

Iceberg 提供了建表的 API,用户可以使用该 API 指定表明、schema、partition 信息等,然后在 Hive catalog 中完成建表。


Delta

我们最后来说 Delta。Delta 的定位是流批一体的 Data Lake 存储层,支持 update/delete/merge。由于出自 Databricks,spark 的所有数据写入方式,包括基于 dataframe 的批式、流式,以及 SQL 的 Insert、Insert Overwrite 等都是支持的(开源的 SQL 写暂不支持,EMR 做了支持)。与 Iceberg 类似,Delta 不强调主键,因此其 update/delete/merge 的实现均是基于 spark 的 join 功能。在数据写入方面,Delta 与 Spark 是强绑定的,这一点 Hudi 是不同的:Hudi 的数据写入不绑定 Spark(可以用 Spark,也可以使用 Hudi 自己的写入工具写入)。

在查询方面,开源 Delta 目前支持 Spark 与 Presto,但是,Spark 是不可或缺的,因为 delta log 的处理需要用到 Spark。这意味着如果要用 Presto 查询 Delta,查询时还要跑一个 Spark 作业。更为蛋疼的是,Presto 查询是基于 SymlinkTextInputFormat。在查询之前,要运行 Spark 作业生成这么个 Symlink 文件。如果表数据是实时更新的,意味着每次在查询之前先要跑一个 SparkSQL,再跑 Presto。这样的话为何不都在 SparkSQL 里搞定呢?这是一个非常蛋疼的设计。为此,EMR 在这方面做了改进,支持了 DeltaInputFormat,用户可以直接使用 Presto 查询 Delta 数据,而不必事先启动一个 Spark 任务。

在查询性能方面,开源的 Delta 几乎没有任何优化。Iceberg 的 hidden partition 且不说,普通的 column 的统计信息也没有。Databricks 对他们引以为傲的 Data Skipping 技术做了保留。不得不说这对于推广 Delta 来说不是件好事。EMR 团队在这方面正在做一些工作,希望能弥补这方面能力的缺失。

Delta 在数据 merge 方面性能不如 Hudi,在查询方面性能不如 Iceberg,是不是意味着 Delta 一无是处了呢?其实不然。Delta 的一大优点就是与 Spark 的整合能力(虽然目前仍不是很完善,但 Spark-3.0 之后会好很多),尤其是其流批一体的设计,配合 multi-hop 的 data pipeline,可以支持分析、Machine learning、CDC 等多种场景。使用灵活、场景支持完善是它相比 Hudi 和 Iceberg 的最大优点。另外,Delta 号称是 Lambda 架构、Kappa 架构的改进版,无需关心流批,无需关心架构。这一点上 Hudi 和 Iceberg 是力所不及的。

delta
image.png


总结

通过上面的分析能够看到,三个引擎的初衷场景并不完全相同,Hudi 为了 incremental 的 upserts,Iceberg 定位于高性能的分析与可靠的数据管理,Delta 定位于流批一体的数据处理。这种场景的不同也造成了三者在设计上的差别。尤其是 Hudi,其设计与另外两个相比差别更为明显。随着时间的发展,三者都在不断补齐自己缺失的能力,可能在将来会彼此趋同,互相侵入对方的领地。当然也有可能各自关注自己专长的场景,筑起自己的优势壁垒,因此最终谁赢谁输还是未知之数。

下表从多个维度对三者进行了总结,需要注意的是此表所列的能力仅代表至 2019 年底。

· Delta Hudi Iceberg
Incremental Ingestion Spark Spark Spark
ACID updates HDFS, S3 (Databricks), OSS HDFS HDFS, S3
Upserts/Delete/Merge/Update Delete/Merge/Update Upserts/Delete No
Streaming sink Yes Yes Yes(not ready?)
Streaming source Yes No No
FileFormats Parquet Avro,Parquet Parquet, ORC
Data Skipping File-Level Max-Min stats + Z-Ordering (Databricks) File-Level Max-Min stats + Bloom Filter File-Level Max-Min Filtering
Concurrency control Optimistic Optimistic Optimistic
Data Validation Yes (Databricks) No Yes
Merge on read No Yes No
Schema Evolution Yes Yes Yes
File I/O Cache Yes (Databricks) No No
Cleanup Manual Automatic No
Compaction Manual Automatic No

注:限于本人水平,文中内容可能有误,也欢迎读者批评指正!


数据湖

图片

数据仓库 VS 数据湖

相较而言,数据湖是较新的技术,拥有不断演变的架构。数据湖存储任何形式(包括结构化和非结构化)和任何格式(包括文本、音频、视频和图像)的原始数据。根据定义,数据湖不会接受数据治理,但专家们一致认为良好的数据管理对预防数据湖转变为数据沼泽不可或缺。数据湖在数据读取期间创建模式。与数据仓库相比,数据湖缺乏结构性,而且更灵活,并且提供了更高的敏捷性。值得一提的是,数据湖非常适合使用机器学习和深度学习来执行各种任务,比如数据挖掘和数据分析,以及提取非结构化数据等。
图片

Apache iceberg:Netflix 数据仓库的基石

Apache iceberg:Netflix 数据仓库的基石

5-year challenges

智能处理引擎

  • CBO,更好的join实现
  • 缓存结果集,物化视图

减少人工维护数据

  • data librarian services 数据图书馆服务
  • declarative instead of imperative 陈述式而不是命令式

Problem Whack-a-mole

1、不安全的操作随处可见: 同时写多个分区,列重命名
2、和对象存储交互有时候会出现很大的问题: eventual consistency to performance problems(最终一致性的性能问题)、output committees can’t fix it
3、无休止的可扩展性挑战。

iceberg

  1. 在单个文件中修改或跳过数据
  2. 当然多个文件也支持这些操作

Hive 表的核心思想是把数据组织成目录树,如上所述。

如果我们需要过滤数据,可以在 where 里面添加分区相关的信息。

带来的问题是如果一张表有很多分区,我们需要使用 HMS(Hive MetaStore)来记录这些分区,同时底层的文件系统(比如 HDFS)仍然需要在每个分区里面记录这些分区数据。

这就导致我们需要在 HMS 和 文件系统里面同时保存一些状态信息;因为缺乏锁机制,所以对上面两个系统进行修改也不能保证原子性。

当然 Hive 这样维护表也不是没有好处。这种设计使得很多引擎(Hive、Spark、Presto、Flink、Pig)都支持读写 Hive 表,同时支持很多第三方工具。简单和透明使得 Hive 表变得不可或缺的。

Iceberg 的目标包括:

1、成为静态数据交换的开放规范,维护一个清晰的格式规范,支持多语言,支持跨项目的需求等。
2、提升扩展性和可靠性。能够在一个节点上运行,也能在集群上运行。所有的修改都是原子性的,串行化隔离。原生支持云对象存储,支持多并发写。
3、修复持续的可用性问题,比如模式演进,分区隐藏,支持时间旅行、回滚等。

Iceberg 主要设计思想:

记录表在所有时间的所有文件,和 Delta Lake 或 Apache Hudi 一样,支持 snapshot,其是表在某个时刻的完整文件列表。每一次写操作都会生成一个新的快照。

读取数据的时候使用当前的快照,Iceberg 使用乐观锁机制来创建新的快照,然后提交。

Iceberg 这么设计的好处是:

  • 所有的修改都是原子性的;
  • 没有耗时的文件系统操作;
  • 快照是索引好的,以便加速读取;
  • CBO metrics 信息是可靠的;
  • 更新支持版本,支持物化视图。

Iceberg 在 Netflix 生产环境维护着数十 PB 的数据,数百万个分区。对大表进行查询能够提供低延迟的响应。

未来工作:1、支持 Spark 向量化以便实现快速的 bulk read,Presto 向量化已经支持。2、行级别的删除,支持 MERGE INTO 等

Storm-runtime

master-slave结构:

  • Nimbus是主节点,负责分发用户代码,指派Supervisor上的worker进程,运行topology的(Spout/Bolt)Task
  • Supervisor是从节点,守护进程. 负责启动和终止worker进程. 通过Storm的配置文件中的 supervisor.slots.ports配置项,可以指定在一个Supervisor上最大允许多少个Slot,每个Slot通过端口号来唯一标识,一个端口号 对应一个Worker进程(如果该Worker进程被启动)。

运行流程

1)户端提交拓扑到nimbus。

2) Nimbus针对该拓扑建立本地的目录根据topology的配置计算task,分配task,在zookeeper上建立assignments节点存储 task和supervisor机器节点中woker的对应关系;

在zookeeper上创建taskbeats节点来监控task的心跳;启动topology。

3) Supervisor去zookeeper上获取分配的tasks,启动多个woker进行,每个woker生成task,一个task一个线程;根据topology 信息初始化建立task之间的连接;Task和Task之间是通过zeroMQ管理的;后整个拓扑运行起来。

内核-队列

在Storm中大量使用Disrupt Queue来解耦Storm内部的消息处理过程。分析这些Queue的分布,是分析Storm运行时态的基础。

在Storm的worker中,最小的执行单元是executor,一个executor只会有一个component。目前一个component只会有一个task。而一个executor会有两个Disruptor Queue, 一个用于接受数据的receive Disrupor Queue和一个用于发送send Disrupor Queue。然后在worker中有一个全局的send Disrupor Queue。分析完队列的分布,在分析topology的运行时状况。

当一个Topology提交到Storm集群后,task被分配到各个wrker中开始执行后,task是怎么执行的。Storm的Task分为两类,一类是消息的源头Spout,它负责在源头产生消息;然后就是Bolt,它是执行单元。但是这两者各自的工作逻辑如下。

我们来分析Spout在运行时的工作状况。Spout的nextTuple用于发送数据,接口注释上说明它不能阻塞,因为它与active, deactive,ack, fail在一个处理线程里面被处理。但是nextTuple被阻塞会有什么副作用列?当nextTuple被阻塞,应用代码中另起线程调用SpoutCollector会有什么后果?下图是SpoutExecutor的执行逻辑。

在SpoutExecutor处理循环里面,第一步做的事情是从receive Disrupor Queue里面消费里面的消息。这里面的就是SpoutExecutor所收到的消息。处理逻辑如下图所示

然后SpoutExecutor会将overflow中的数据再次发送。overflow是用于接受SpoutCollector.emit()无发及时发送的数据的,具体SpoutCollector发送数据的逻辑见下文分析。但是这里在发送overflow的数据时,与SpoutCollector.emit()有个区别,就是数据还是无法被正常发送时,数据会丢弃,也就是不会被再次写入overflow中。当overflow中没有数据,以及pending中的数据量小于TOPOLOGY_MAX_SPOUT_PENDING时,判断topology的状态。当topology不是deactive状态时,如果topology有触发active命令,会调用spout的active接口,然后调用spout的nextTuple接口。否则调用deactive接口。所以这里当Spout的nextTuple被阻塞时,spout没办法处理acker回报的消息,回报的消息会阻塞在executor的receive DisruptorQueue中,当receive DisruptorQueue塞满后,是会阻塞在对应的网络处理模块中,storm中经典的是zeroMQ,老版本的zeroMQ是没有设置水位,这样会大量堆积到内存中,因为zeroMQ是C++的,占用的是堆外内存,JVM无法管理,最坏就是把机器的内存耗光。如果这个是另起线程调用SpoutCollector的emit,当emit数据速度过快,会导致overflow中堆积数据,导致worker内存消耗。

上面讲到SpoutCollector在emit数据会写入overflow,那什么情况下会写入overflow。SpoutCollector.emit是否真的就把数据发送到了网络。下面是SpoutCollector的处理逻辑。

Spout当调用SpoutCollect.emit()发送时,首先的逻辑是根据根据消息的STREAM ID,和对应的values,得到目的端task id。当topology开启acker机制时,会生成一个随机的rootId,否则使用一个默认值作为rootId。然后为消息生成一个随机的messageId,由rootId和messageId组成对应的tuple Id。这里tuple Id是有rootId为key,messageId为value的一个HashMap。这里采用HashMap的作用会在下面Spout的ack机制是做说明。在将生成的tuple Id和用会的values组成一个tuple。并判断overflower是否有数据,让overflow中有数据时,数据会直接放入overflow中,而不会放入executor 的send Disraptor Queue中。当overflow为空时,会将tuple放入send Disraptor Queue中。当捕获Disraptor Queue的InsufficientCapacityException时,数据就放入了overflow中。然后处理ack。当开启了ack机制,会在pending中加入对应的tuple数据,然后项acker发送ack init消息。

根据上面的发送过程,数据emit只是被写入了executor的send Disraptor Queue。而数据在Spout端的丢失多是数据在overflow中被SpoutExecutor在次处理时,send Disraptor Queue满导致。

关于pending,首先它是一个RotatingMap。它通过定时旋转,以达到定时器的目的。在SpoutExecutor的RotatingMap中有两个桶,然后executor有个定时线程,会按照用户设定的message timeout second,定期向executor的receive DisruptorQueue写入SYSTEM_TICK_STREAM_ID消息,然后SpoutExecutor处理SYSTEM_TICK_STREAM_ID消息时就旋转RotatingMap,当被清除的桶中有数据时,被清除的桶中的数据会调用fail接口,通知业务逻辑fail。这就是当超时间设置比业务逻辑短时,导致数据重复的原因。还有一种是acker消息丢失,导致数据重复。由于RotatingMap的底层是HashMap,中间没有锁,overflow是个LinkedList,所以SpoutCollector和SpoutExecutor不能在两个线程中并发 处理。

BoltExecutor的处理逻辑非常简单,就是消费receive Disrupor Queue中数据,然后调用bolt的excue。但是这里也是单线程处理的,阻塞或者处理速度不匹配,就会导致数据在Disrupor Queue或者网络模块中堆积,其中使用zeroMQ的副作用最大。

BoltCollector的emit与SpoutCollector的emit处理相比,首先是少了overflow承接无法发送的数据,会直接丢弃。其次是没有pending和acker消息的发送。

接下来分析一下Storm的Acker机制。Storm的Ack机制在Storm刚刚开源时被大书特书,处理原理也确实非常的精彩。由于这里ID都是随机数,所以这里不会在讨论随机数的唯一问题。

如上图所示,spout发送t1给bolt a, bolt a 在t1的基础上生成t2,t3,t4给bolt b,bolt b ack所收到的数据。下面来追踪整个id变化的过程。

Spout 发送t1, message id为<r_1, m_1>,发送ack

Acker收到ack init, map中缓存<r_1,m_1>

Bolt a 收到t1, messageId为<r_1,m_1>,生成t2,t3, t4

Bolt a 发送t2, anchor t1, t2,messageId为<r_1,m_2>, 更新t1的ackVal为m_2

Bolt a 发送t3, anchor t1, t3,messageId为<r_1,m_3>, 更新t1的ackVal为m_2^m_3

Bolt a 发送t4, anchor t1, t4,messageId为<r_1,m_4>, 更新t1的ackVal为m_2^m_3^m_4

Bolt a ack t1, 向acker发送ack消息<r_1, m_1^ m_2^m_3^m_4>

Acker 收到bolt a的ack消息,更新缓存为<r_1, m_1^ m_1^ m_2^m_3^m_4>即<r_1, m_2^m_3^m_4>

Bolt b ack t2, 向acker发送ack消息<r_1, m_2>

Acker 收到bolt b的ack消息,更新缓存为<r_1, m_2^ m_2^m_3^m_4>即<r_1, m_3^m_4>

Bolt b ack t3, 向acker发送ack消息<r_1, m_3>

Acker 收到bolt b的ack消息,更新缓存为<r_1, m_3^m_3^m_4>即<r_1, m_4>

Bolt b ack t4, 向acker发送ack消息<r_1, m_4>

Acker 收到bolt b的ack消息,更新缓存为<r_1, m_4^m_4>即<r_1, 0>

messageId确认完毕,向Spout发送ack消息。当消息没有被ack,会一直在spout的pending队列中,知道被ack或者超时。

它基本上使用两个long值就跟踪了一个消息在整个流中的处理过程。

【参考文献】


  1. Apache storm内核原理

kafka

ACK机制

Kafka producer有三种ack机制 初始化producer时在config中进行配置

ACK 同步 延迟
0 producer不等待broker同步完成就发送下一条(批)信息 低的延迟最弱的持久性,当服务器发生故障时,就很可能发生数据丢失。例如leader已经死亡,producer不知情,还会继续发送消息broker接收不到数据就会数据丢失
1 producer要等待leader成功收到数据并得到确认,才发送下一条message 较好的持久性较低的延迟性:Partition的Leader死亡,follwer尚未复制,数据就会丢失
-1 producer得到follwer确认,才发送下一条数据 持久性最好,延时性最差

三种机制性能递减,可靠性递增

GC调优

TODO

内容:

  1. 教程
  2. 压测+调优
  3. 实际样例

tip:

  • collector
  • gc logs
  • gc viewer
  • jmeter
  • 压测与调优

JVM内存结构

jvm整体架构图文详解

Java8 语言规范
Java8 JVM规范
Java8 JVM规范-内存结构

运行时数据区

程序计数器

虚拟机栈JVM Stacks

栈帧分为哪些快?每块又保存什么内容?

方法区

非堆

  • JDK6 perm区
  • JDK7 perm区
  • JDK8 metaspace

常量池Run-Time Constant Pool (方法区中)

本地方法栈Native Method Stacks

JVM的内存结构

  • CCS:只有启用了短指针的时候,才存在
  • CodeCache:只有启用了JIT和有JNI调用Native代码的时候,才存在
    • -Xcomp:JIT完全编译执行
    • -Xint完全解释执行
    • -Xmixed编译和解释混合

非堆区

标准参数

1
2
3
4
-version -showversion
-help
-cp -classpath
-server -client

-X

-Xint: 解释执行模式
-Xcomp: 编译执行模式, 第一次使用就编译成本地代码, 编译结果保存在metaspace的code cache空间
-Xmixed: 混合执行模式, JVM决定是否编译成本地代码

1
2
3
4
5
6
7
8
9
> java -Xint -version
openjdk version "1.8.0_232"
OpenJDK Runtime Environment (build 1.8.0_232-b09)
OpenJDK 64-Bit Server VM (build 25.232-b09, interpreted mode)

> java -Xcomp -version
openjdk version "1.8.0_232"
OpenJDK Runtime Environment (build 1.8.0_232-b09)
OpenJDK 64-Bit Server VM (build 25.232-b09, compiled mode)

-XX

参数 作用
-Xms 最小堆内存
-Xmx 最大堆内存
-XX:NewSize 新生代大小
-XX:MaxNewSize 最大新生代大小
-XX:NewRatio new区和old区的比例
-XX:SurvivorRatio eden区与survivor区大小比例
-XX:MetaspaceSize Metaspace大小
-XX:MaxMetaspaceSize Metaspace最大大小
-XX:+UseCompressedClassPointers 压缩类指针
-XX:CompressedClassSpaceSize 压缩类空间(CCS)的大小,默认1G
-XX:InitialCodeCacheSize code cache的初始大小
-XX:ReservedCodeCacheSize code cache的最大的大小
-XX:PretenureSizeThreshold 大对象直接进入老年代,大对象的大小阈值
-XX:MaxTenuringThreshold 长期存活的对象进入老年代,晋升年龄阈值
-XX:+PrintTrnuringDistuibution youngGC时打印年龄分布情况
-XX:TargetSurvivorRatio survivor垃圾回收存活的比例,超过值将直接晋升

垃圾回收算法

gc调优官方指南

GC Root:

  • 类加载器:由类加载器生成的对象,都持有指针
  • Thread:线程运行会持有很多对象
  • 虚拟机栈的本地变量表
  • static成员
  • 常量引用
  • 本地方法栈的变量

引用计数

缺点:
无法处理循环引用

标记清除

先标记需要回收的对象,在统一回收所有对象

缺点:

效率不高:标记和清除两个过程效率都不高;碎片:导致提前GC

复制

内存划分为大小相同的两块,每次只使用其中一块,一块用完复制存活的对象到另一块,然后再把已使用的内存空间一次清理掉

缺点:

使用简单,效率高,空间利用率不高

标记整理

先标记需要回收的对象,让所有存活的对象都向一端移动,然后清理掉端边界外的内存

缺点:
无内存碎片,比较耗时

分带垃圾回收

young区朝生夕死,生命周期端,用复制算法:效率高
Old区生命周期长,用标记清除或标记整理

  • 对象优先分配在eden区
  • 大对象直接进入老年代:-XX:PretenureSizeThreshold
  • 长期存活的对象进入老年代:
    • -XX:MaxTenuringThreshold: 晋升年龄代数阈值
    • -XX:+PrintTenuringDistribution:ygc打印存活对象的分布情况
    • -XX:TargetSurvivorRatio:Survivor区存活对象比例,动态调整,取存活对象的平均值与晋升年龄阈值间的最小值

垃圾收集器

枚举根节点,做可达性分析
根节点: 类加载器、Thread、虚拟机栈的本地变量表、static成员、常量引用、本地方法栈的变量

  • 串行收集器Serial: Serial、 Serial old
  • 并行收集器Parallel: Parallel Scavenge、Parallel old,吞吐量优先
  • 并发收集器Concurrent: CMS、G1,停顿时间优先

并行 vs 并发

并行(Parallel): 多条垃圾收集线程并行工作,但此时用户线程仍然处于等待状态。适合科学计算、后台处理等弱交互的场景

并发(Concurrent): 用户线程和垃圾收集线程同时执行(但不一定是并行的,可能会交替执行),垃圾收集线程在执行的时候不会停顿用户程序的运行。适合对响应时间有要求的场景,如web。

停顿时间 vs 吞吐量

停顿时间:垃圾收集器做垃圾回收中断应用执行的时间。-XX:MaxGCPauseMillis

吞吐量:花在垃圾收集的时间和花在应用时间的占比。 -XX:GCTimeRatio=<n>, 垃圾收集时间占: 1/(1+n)

串行收集器

-XX:+UseSerialGC
-XX:+UseSerialOldGC

采用串行收集器,默认old区采用串行收集器

并行收集器 ParallelCollector

吞吐量优先

1
2
3
4
-XX:+UseParallelGC
-XX:+UseParallelOldGC

Server模式下的默认收集器
1
2
3
4
-XX:ParallelGCThreads=<N> 多少个GC线程

CPU>8 N=5/8
CPU<8 N=CPU

并发收集器

响应时间优先

CMS: -XX:+UseConcMarkSweepGC -XX:+UseParNewGC
G1: -XX:+UseG1GC

如何选择垃圾收集器

如何选择垃圾收集器

  • 优先调整堆的大小让服务器自己选择
  • 如果内存小于100M,使用串行收集器
  • 如果是单核,并且没有停顿时间的要求,串行或者jvm自己选
  • 如果允许停顿时间超过1s,选择并行或者jvm自己选
  • 如果响应时间最重要,并且不能超过1s,则使用并发收集器
young Tenured JVM options
Serial Serial -XX:+UseSerialGC
Parallel Scavenge Serial -XX:+UseParallelGC -XX:-UseParallelOldGC
Parallel Scavenge Parallel Old -XX:+UseParallelGC -XX:+UseParallelOldGC
Parallel New或Serial CMS -XX:+UseParNewGC -XX:+UseConcMarkSweepGC
G1 G1 -XX:+UseG1GC

垃圾回收器从线程运行情况分类有三种

串行回收: Serial回收器,单线程回收,全程stw;
并行回收: 名称以Parallel开头的回收器,多线程回收,全程stw;
并发回收: cms与G1,多线程分阶段回收,只有某阶段会stw;

并行收集器 Parallel Collector

暂停应用程序,开启多个垃圾收集线程开始垃圾回收

  • -XX:+UseParallelGC 手动开启,Server默认开启
  • -XX:ParallelGCThreads=<N>多少个GC线程
    • CPU>8 N=5/8
    • CPU<8 N=CPU

查找使用ParallelGC的进程
jps -v | grep -v grep | awk '{print $1}' | xargs -L 1 -t jinfo -flag UseParallelGC

Parallel Collector Ergonomics自适应

  • -XX:MaxGCPauseMillis=<N>:最大停顿时间
  • -XX:GCTimeRatio=<N>: GC时间占比,代表吞吐量
  • -Xmx<N>: 堆最大大小

优先满足停顿时间,再满足吞吐量的要求,最后再调整满足堆最大大小

动态调整每个分区的大小

动态内存调整

  • -XX:YoungGenerationSizeIncrement=<Y> 年轻代大小调整增量,默认值20%
  • -XX:TenuredGenerationSizeIncrement=<T> 老年代大小调整增量,默认值
  • -XX:AdaptiveSizeDecrementScaleFactor=<D> 减少增量,默认值4%

CMS

1
java -XX:+UseConcMarkSweepGC  -jar -server console.jar
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
jps -l | grep buried | awk '{print $1}' | xargs -L 1 -t /usr/local/soft/jdk1.8.0_191/bin/jinfo  -flags
/usr/local/soft/jdk1.8.0_191/bin/jinfo -flags 4893
Attaching to process ID 4893, please wait...
Debugger attached successfully.
Server compiler detected.
JVM version is 25.191-b12
Non-default VM flags:
-XX:CICompilerCount=3
-XX:InitialHeapSize=524288000
-XX:MaxHeapSize=8363442176
-XX:MaxNewSize=348913664
-XX:MaxTenuringThreshold=6
-XX:MinHeapDeltaBytes=196608
-XX:NewSize=174718976
-XX:OldPLABSize=16
-XX:OldSize=349569024
-XX:+UseCompressedClassPointers
-XX:+UseCompressedOops
-XX:+UseConcMarkSweepGC
-XX:+UseFastUnorderedTimeStamps
-XX:+UseParNewGC
Command line: -XX:+UseConcMarkSweepGC
1
2
3
jps -l | grep buried | awk '{print $1}' | xargs -L 1 -t /usr/local/soft/jdk1.8.0_191/bin/jinfo  -flag CMSInitiatingOccupancyFraction
/usr/local/soft/jdk1.8.0_191/bin/jinfo -flag CMSInitiatingOccupancyFraction 4893
-XX:CMSInitiatingOccupancyFraction=-1

cms是一种预处理垃圾回收器,它不能等到old内存用尽时回收,需要在内存用尽前,完成回收操作,否则会导致并发回收失败;所以cms垃圾回收器开始执行回收操作,有一个触发阈值,默认是老年代或永久带达到92%

  • 并发收集
  • 低停顿 低延迟
  • 老年代收集器

CMS垃圾收集过程

  1. CMS inital mark: 初始标记Root STW
  2. CMS concurrent mark:并发标记
  3. CMS-concurrent-preclean: 并发预清理
  4. CMS remark: 重新标记 STW
  5. CMS concurrent sweep:并发清除
  6. CMS-concurrent-reset:并发重置

缺点

  • 低停顿 低延迟
  • CPU敏感
  • 浮动垃圾:边运行应用程序,边回收
  • 空间碎片

调优参数

参数 备注
-XX:ConcGCThreads 并发的GC线程数
-XX:+UseCMSCompactAtFullCollection FullGC之后做压缩
-XX:CMSFullGCsBeforeCompaction 多少次FullGC之后压缩一次
-XX:CMSInitiatingOccupancyFraction 触发FullGC 92%
-XX:+UseCMSInitiatingOccupancyOnly 是否动态调
-XX:+CMSScavengeBeforeRemark FullGC之前先做YGC
-XX:+CMSClassUnloadingEnabled 启用回收Perm区

G1

大内存(大于6G),优先延迟(小于0.5s)

H区:大对象,如果对象超过了region的一半大小

Region

SATB:snapshot-at-the-beginning, 通过Root tracing得到的,GC开始时候存活对象的快照。垃圾回收以此为基础回收

RSet:记录了其他Region中的对象引用本Region中对象的关系,属于points-into结构(谁引用了我的对象)

YoungGC

  • 新独享进入Eden区
  • 存活对象拷贝到s区
  • 存活时间达到年龄阈值时,对象晋升到old区

mixedGC

没有full gc

  • 不是FullGC,回收所有的Young和部分Old
  • global concurrent marking

global concurrent marking

  1. Initial marking phase:标记GC Root ,STW
  2. Root region scanning phase:标记存活Region
  3. Concurrent marking phase:标记存活的对象
  4. Remark phase:重新标记 STW
  5. Cleanup phase:部分STW

MixedGC时机

  • InitiatingHeapOccupancyPercent: 堆占有率达到这个数值则触发global concurrent marking,默认45%
  • G1HeapWastePercent:在gloabl concurrent marking结束之后,可以知道区有多少空间要被回收,在每次YGC之后和再次发生MixedGC之前,会检查垃圾占比是否达到此参数,只有达到了,下次才会发生MixedGC
  • G1MixedGCLiveThresholdPercent: Old区的region被回收时候的存活对象占比
  • G1MixedGCCountTarget:一次global concurrent marking之后,最多执行MixedGC的次数

调优最佳实践

可视化GC日志分析工具

吞吐量与延迟时间的权衡

Tomcat调优实例

CMS垃圾回收器详解

转载自Java并发编程之原子性-Atomic详解

Atomic

atomic包

JUC中的Atomic包详解:

Atomic包中提供了很多Atomicxxx的类:

img

它们都是CAS(compareAndSwap)来实现原子性。

AtomicInteger样例

先写一个简单示例如下:

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
29
30
31
32
33
34
35
@Slf4j
public class AtomicExample1 {

// 请求总数
public static int clientTotal = 5000;
// 同时并发执行的线程数
public static int threadTotal = 200;

public static AtomicInteger count = new AtomicInteger(0);

public static void main(String[] args) throws Exception {
ExecutorService executorService = Executors.newCachedThreadPool();
final Semaphore semaphore = new Semaphore(threadTotal);
final CountDownLatch countDownLatch = new CountDownLatch(clientTotal);
for (int i = 0; i < clientTotal ; i++) {
executorService.execute(() -> {
try {
semaphore.acquire();
add();
semaphore.release();
} catch (Exception e) {
log.error("exception", e);
}
countDownLatch.countDown();
});
}
countDownLatch.await();
executorService.shutdown();
log.info("count:{}", count.get());
}

private static void add() {
count.incrementAndGet();
}
}

可以发下每次的运行结果总是我们想要的预期结果5000。 说明该计数方法是线程安全的。

AtomicInteger实现原理

我们查看下count.incrementAndGet()方法,它的第一个参数为对象本身,第二个参数为valueOffset是用来记录value本身在内存的编译地址的,这个记录,也主要是为了在更新操作在内存中找到value的位置,方便比较,第三个参数为常量1。:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
public class AtomicInteger extends Number implements java.io.Serializable {
private static final long serialVersionUID = 6214790243416807050L;

// setup to use Unsafe.compareAndSwapInt for updates
private static final Unsafe unsafe = Unsafe.getUnsafe();
private static final long valueOffset;

static {
try {
valueOffset = unsafe.objectFieldOffset
(AtomicInteger.class.getDeclaredField("value"));
} catch (Exception ex) { throw new Error(ex); }
}
private volatile int value;
... 此处省略多个方法...
/**
* Atomically increments by one the current value.
*
* @return the updated value
*/
public final int incrementAndGet() {
return unsafe.getAndAddInt(this, valueOffset, 1) + 1;
}
}

AtomicInteger源码里使用了一个Unsafe的类,它提供了一个getAndAddInt的方法,我们继续点看查看它的源码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
public final class Unsafe {
private static final Unsafe theUnsafe;

....此处省略很多方法及成员变量....


public final int getAndAddInt(Object var1, long var2, int var4) {
int var5;
do {
var5 = this.getIntVolatile(var1, var2);
} while(!this.compareAndSwapInt(var1, var2, var5, var5 + var4));

return var5;
}

public final native boolean compareAndSwapInt(Object var1, long var2, int var4, int var5);

public native int getIntVolatile(Object var1, long var2);
}

可以看到这里使用了一个do while语句来做主体实现的。而在while语句里它的核心是调用了一个compareAndSwapInt()的方法。它是一个native方法,它是一个底层的方法,不是使用Java来实现的。

假设我们要执行0+1=0的操作,下面是单线程情况下各参数的值:

img img 更新后:

img

compareAndSwapInt()方法的第一个参数(var1)是当前的对象,就是代码示例中的count。此时它的值为0(期望值)。 第二个值(var2)是传递的valueOffset值,它的值为12。 第三个参数(var4)就为常量1。方法中的变量参数(var5)是根据参数一和参数二valueOffset,调用底层getIntVolatile方法得到的值,此时它的值为0 。 compareAndSwapInt()想要达到的目标是对于count这个对象,如果当前的期望值var1里的value跟底层的返回的值(var5)相同的话,那么把它更新成var5+var4这个值。 不同的话重新循环取期望值(var5)直至当前值与期望值相同才做更新。compareAndSwap方法的核心也就是我们通常所说的CAS。

Atomic包下其他的类如AtomicLong等的实现原理基本与上述一样。

AtomicInteger的代码

image.png

他的值是存在一个volatile的int里面。volatile只能保证这个变量的可见性。不能保证他的原子性。

可以看看getAndIncrement这个类似i++的函数,可以发现,是调用了UnSafe中的getAndAddInt。

image.png

UnSafe

UnSafe是何方神圣?UnSafe提供了java可以直接操作底层的能力。

进一步,我们可以发现实现方式:

image.png

如何保证原子性:自旋 + CAS(乐观锁)。在这个过程中,通过compareAndSwapInt比较更新value值,如果更新失败,重新获取旧值,然后更新。

优缺点

CAS相对于其他锁,不会进行内核态操作,有着一些性能的提升。但同时引入自旋,当锁竞争较大的时候,自旋次数会增多。cpu资源会消耗很高

换句话说,CAS+自旋适合使用在低并发有同步数据的应用场景。

Java 8做出的改进和努力

在Java 8中引入了4个新的计数器类型,LongAdderLongAccumulatorDoubleAdderDoubleAccumulator。他们都是继承于Striped64

在LongAdder 与AtomicLong有什么区别?

这里再介绍下LongAdder这个类,通过上述的分析,我们已经知道了AtomicLong使用CAS: 在一个死循环内不断尝试修改目标值直到修改成功。如果在竞争不激烈的情况下,它修改成功概率很高。反之,如果在竞争激烈的情况下,修改失败的概率会很高,它就会进行多次的循环尝试,因此性能会受到一些影响。 对于普通类型的long和double变量,jvm允许将64位的读操作或写操作拆成两个32位的操作。

LongAdder的核心思想是将热点数据分离,它可以将AtomicLong内部核心数据value分离成一个数组,每个线程访问时通过hash等算法映射到其中一个数字进行计数。而最终的计数结果则为这个数组的求和累加,其中热点数据value,它会被分离成多个单元的cell,每个cell独自维护内部的值,当前对象的实际值由所有的cell累计合成。这样,热点就进行了有效的分离,提高了并行度。LongAdder相当于在AtomicLong的基础上将单点的更新压力分散到各个节点上,在低并发的时候对base的直接更新可以很好的保障跟Atomic的性能基本一致。而在高并发的时候,通过分散提高了性能。但是如果在统计的时候有并发更新,可能会导致统计的数据有误差。

在实际高并发计数的时候,可以优先使用LongAdder。在低并行度或者需要准确数值的时候可以优先使用AtomicLong,这样反而效率更高。

下面简单的演示下Atomic包下AtomicReference简单的用法:

1
2
3
4
5
6
7
8
9
10
11
@Slf4j
public class AtomicExample4 {

private static AtomicReference<Integer> count = new AtomicReference<>(0);

public static void main(String[] args) {
count.compareAndSet(0, 2);
count.compareAndSet(0, 1);
log.info("count:{}", count.get());
}
}

compareAndSet()分别传入的是预期值跟更新值,只有当预期值跟当前值相等时,才会将值更新为更新值;

上面的第一个方法可以将值更新为2,而第二个步中无法将值更新为1。

Atomic*遇到的问题是,只能运用于低并发场景。因此LongAddr在这基础上引入了分段锁的概念。可以参考《JDK8系列之LongAdder解析》一起看看做了什么。

大概就是当竞争不激烈的时候,所有线程都是通过CAS对同一个变量(Base)进行修改,当竞争激烈的时候,会将根据当前线程哈希到对于Cell上进行修改(多段锁)。

image.png

可以看到大概实现原理是:通过CAS乐观锁保证原子性,通过自旋保证当次修改的最终修改成功,通过**降低锁粒度(多段锁)**增加并发性能。

AtomicIntegerFieldUpdater

下面简单介绍下AtomicIntegerFieldUpdater 用法(利用原子性去更新某个类的实例):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@Slf4j
public class AtomicExample5 {

private static AtomicIntegerFieldUpdater<AtomicExample5> updater =
AtomicIntegerFieldUpdater.newUpdater(AtomicExample5.class, "count");

@Getter
private volatile int count = 100;

public static void main(String[] args) {

AtomicExample5 example5 = new AtomicExample5();

if (updater.compareAndSet(example5, 100, 120)) {
log.info("update success 1, {}", example5.getCount());
}

if (updater.compareAndSet(example5, 100, 120)) {
log.info("update success 2, {}", example5.getCount());
} else {
log.info("update failed, {}", example5.getCount());
}
}
}

它可以更新某个类中指定成员变量的值。注意:修改的成员变量需要用volatile关键字来修饰,并且不能是static描述的字段。

AtomicStampReference

AtomicStampReference 这个类它的核心是要解决CAS的ABA问题(CAS操作的时候,其他线程将变量的值A改成了B,接着又改回了A,等线程使用期望值A与当前变量进行比较的时候,发现A变量没有变,于是CAS就将A值进行了交换操作。实际上该值已经被其他线程改变过)。ABA问题的解决思路就是每次变量变更的时候,就将版本号加一。看一下它的一个核心方法compareAndSet():

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
29
public class AtomicStampedReference<V> {

private static class Pair<T> {
final T reference;
final int stamp;
private Pair(T reference, int stamp) {
this.reference = reference;
this.stamp = stamp;
}
static <T> Pair<T> of(T reference, int stamp) {
return new Pair<T>(reference, stamp);
}
}

... 此处省略多个方法 ....

public boolean compareAndSet(V expectedReference,
V newReference,
int expectedStamp,
int newStamp) {
Pair<V> current = pair;
return
expectedReference == current.reference &&
expectedStamp == current.stamp &&
((newReference == current.reference &&
newStamp == current.stamp) ||
casPair(current, Pair.of(newReference, newStamp)));
}
}

可以看到它多了一个stamp的比较,stamp的值是由每次更新的时候进行维护的。

AtomicLongArray

再介绍下 AtomicLongArray , 它维护了一个数组。在该数组下,我们可以选择性的已原子性操作更新某个索引对应的值。