0%

Kettle中Java节点生成

Pentaho – Kettle | TechBI - Business Intelligence

一、问题背景

现有的Transformer无法满足需求,需要购买第三方插件、开发插件、自定义Java类等办法,自定义Java类,因其代价小且门槛较低而成为最为常用的定制方法。

二、原理剖析

Kettle转换的核心对象/脚本类别,属于典型的需要编程基础才能掌控的Transformer类型。而Java代码Transformer,适用于熟悉Java语言的开发人员,用好这个Transformer,需要对类、接口、多线程等语言相关知识有所掌握,并且需要对Kettle的基础框架有所理解。

Kettle转换的执行,有以下三个核心的生命周期节点:

1、初始化

Kettle转换在执行前,会有一个各Transformer的初始化动作,为Transformer执行前的准备工作创造机会。为提高初始化的性能,Kettle为每个Transformer启用一个初始化线程,从而并行完成所有Transformer的初始化。初始化的主要内容就是调用一次Transformer的以下方法:

1
public boolean init( StepMetaInterface meta, StepDataInterfacedata)

此方法包含两个参数。其中,meta为元数据,data为数据。如果返回true,那么代表初始化成功,否则代表初始化失败。任何一个Transformer初始化失败,都会导致整个转换停止执行(在停止前,会调用每一个转换的资源释放方法dispose)。

2、执行

执行阶段是每一个Transformer实现特定价值的时候。为提高效率,Kettle为每一个Transformer单独启动一个工作线程来执行任务。Java程序员都了解,线程的核心代码是覆盖run方法。为简化起见,我将不重要的代码删除,得到工作线程run方法核心代码:

img

可以看出,线程一直在执行Transformer的processRow方法,直到出现以下情况之一:

· processRow方法返回false

· isStopped方法返回true

· processRow方法执行过程中出现异常,

其中,第一种情况代表工作已经正常完成;第二种情况,代表Transformer被强制停止;第三种情况,代表执行过程中出现错误,Kettle将调用stopAll方法,从而导致整个转换的所有工作线程停止执行。

执行方法的声明如下:

1
public boolean processRow( StepMetaInterface meta,StepDataInterfacedata ) throws KettleException;

每一个Transformer,都会在processRow方法中各显神通。一般的过程是,从输入行集中拿出一行,进行特定处理,然后将新的行放入输出行集中。从输入行集中取数据可以调用getRow方法。如果getRow方法返回值不为null,则Transformer应将该行数据进行处理,并调用putRow方法将处理结果存入输出行集,然后返回true,以继续为下一行输入数据处理提供机会。如果getRow方法返回null,代表输入行集已经处理完毕,这时可以调用setOutputDone,标识本Transformer执行完毕,并返回false,以结束本工作线程的执行。

3、资源释放

从上述工作线程的核心代码可以看出,不管工作线程是正常执行完毕还是异常执行完毕,最终会调用dispose方法。该方法声明如下:

1
public void dispose( StepMetaInterface meta, StepDataInterfacedata);

Transformer应该在需要时覆盖此方法,并释放相关资源。

了解上述原理后,撰写Java类Transformer中的代码时将胸有成竹。综上所述,一般情况下重写processRow方法即可满足需求,如果用到了一些重量级的资源,最好在init方法中初始化,并在dispose方法中释放。

由于Kettle使用Janino框架为自定义Java转换Transformer类动态定义了类名,并指定父类为TransformClassBase,所以在撰写代码时,只需要提供类的内容即可,无需class声明。

既然自动建立了父类,那么父类的成员、方法都可以在代码中重用。父类常用的成员包括以下三个实例:

parent:代表容器对象

meta:代表容器元数据对象

data:代表容器数据对象

常用的方法包括:

getRow:从输入行集中取一行数据

putRow:存银行数据到输出行集

stopAll:停止所有工作线程

setOutputDone:标记本Transformer工作完成

logBasic:输出基本日志

logError:输出错误日志

getInputRowMeta:得到输入行的元数据

createOutputRow:创建一个输出行数据

其实,常用的方法(如下图1所示),基本上都在Transformer属性对话框左侧Code Snippits中。一般情况下,可以双击其中的Main节点,从processRow方法的重写开始,需要其他代码时,在左侧找到对应代码块,双击即可加入。

img

三、案例分享

本文使用一个Kettle集成JMS的案例来进行实战演练。假设需要两个转换:一个转换名为Send,实现从文本文件输入流读取数据,并发送到ActiveMQ的队列;另外一个转换名为Receive,实现从队列读取数据,并发送到文本文件输出流。两个转换截图如下:

img

由于两个转换中用到的文本文件输入、输出都非常简单,这里只做简单描述。Send转换中,S01读取本地文本文件,包含两个String类型的字段ID、MENU_NAME。Receive转换中,S02输出文件到转换所在目录,仅包含一个名为MENU_NAME的String类型字段。

下文着重描述两个Java代码Transformer。第一个Transformer是Send转换中的S02,其主要代码注释如下:

img

第二个Transformer是Receive转换中的S01Transformer。主要代码注释如下:

img

注意,由于本文使用了ActiveMQ作为JMS服务器,所以为保证实例能够正常运行,需要自行下载服务器安装程序,并将对应jar文件拷贝到Kettle的lib目录下(本例使用activemq-all-5.8.0.jar)。代码中,需要的import指令如下:

1
2
3
4
5
6
7
8
9
10
import java.util. * ;
import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.DeliveryMode;
import javax.jms.Destination;
import javax.jms.MapMessage;
import javax.jms.MessageProducer;
import javax.jms.Session;
import org.apache.activemq.ActiveMQConnection;
import org.apache.activemq.ActiveMQConnectionFactory;

四、总结

本文在具备程序员背景知识的数据工程师在运用Kettle进行定制开发时,可以参考。

Flink 实时规则引擎实现与应用

Rule engine

欺诈检测
实时标签工厂
ABT
CEP有什么问题?

Frund detected system

秒级规则生效
自定义分发,运行时动态切换
自定义窗口
实时多维聚合

样例:

  1. 欺诈检测: 同一个用户连续往两一个账户三天内连续支付超过
  2. 实时标签工厂: 标签-用户近三十天访问次数、近三十天访问某个页面的次数
  3. ABT

实时用户画像

场景

企业面对不断增加的海量信息,其信息筛选和处理效率低下的困扰与日俱增。由于用户营销不够细化,企业App 中许多不合时宜或不合偏好的消息推送很大程度上影响了用户体验,甚至引发了用户流失

玖富实时反欺诈

业务场景

  • 低延迟
  • 大数据
  • 多维度

系统获取用户产生数据最简单有效的方法就是流水式数据,单个数据包里包含了发生时间点的各个维度的所有信息量,这种场景的特性之一就是数据高并发,因此对时效要求比较高的数据分析来说是一个非常巨大的挑战

在哪个场景下,实时宽表都是一个门槛

  • 高并发

  • 规则实时下发

数据加工提速

在大量数据中做快速预查,利用Flink并发能力进行数据覆盖,最后在缓存里命中结果,从而不必重新进行网络I/O 查询、等待返回的过程。经过部分计算框架升级,最终系统实现了p99 延迟由1s 降为100ms 的优化

Flink fraud detect

https://github.com/afedulov/fraud-detection-demo

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
DataStream<Transaction> transactions = getTransactionsStream(env);

DataStream<Rule> rulesUpdateStream =;
BroadcastStream<Rule> rulesStream = rulesUpdateStream.broadcast(Descriptors.rulesDescriptor);

DataStream<Alert> alerts =
transactions
.connect(rulesStream)
.process(new DynamicKeyFunction())
.uid("DynamicKeyFunction")
.name("Dynamic Partitioning Function")
.keyBy((keyed) -> keyed.getKey())
.connect(rulesStream)
.process(new DynamicAlertFunction())
.uid("DynamicAlertFunction")
.name("Dynamic Rule Evaluation Function");


metric

1
2
ruleCounterGauge = new RuleCounterGauge();
getRuntimeContext().getMetricGroup().gauge("numberOfActiveRules", ruleCounterGauge);

在大规模数据处理的分布式系统中,如何保障数据的高可用、数据的一致性和幂等性(exactly once)是系统的一大难题!
使用廉价机器构建集群成为大数据平台的标配,故障恢复和容错(failover recovery)机制成为防止消息丢失和快速恢复服务必不可少的组成部分; 在通用的大数据架构中,也是保障数据高可用、一致性和幂等性的基础。

分布式系统故障恢复需要解决的问题:

  1. 吞吐量大的场景,当出现了失败时,需要保证失败的数据可以重放,且状态可恢复。
  2. 某一个时刻,多个计算节点存在处理速度不一致的问题,一条数据可能经过多个计算节点才能完成计算。如何保存多个计算节点的状态,且保证数据对齐?
  3. 如何保证数据经过各个计算节点的顺序性和不重复?
  4. 如何回放数据?

虽然大规模分布式系统的侧重各不相同,但是failover recovery的机制却如出一辙, 主要有ack模式、异步checkpoint模式、CL模式、补偿模式等,接下来就这四种模式分别总结。

Ack模式

在分布式系统中,为了确保一条(批)数据被正确处理,且当出现任何故障,保障数据不丢。ack机制是最简单的方式,每一条(批)数据正确处理完后,发送一条确认标示。

TCP ack

在TCP握手时,当收到客户端发起的握手报文时, 返回一个Acknowledgement Number, 标示客户端的请求已经收到,并返回给客户端。
而后客户端返回Acknowledgement Number以确保服务端请求正确接受。
当网络抖动或者服务器和客户端故障时,报文可能丢失。此时就要依赖于与ack机制相配合的faiover-recovery机制。

如果服务器没有收到客户端的最终ACK确认报文,会一直处于SYN_RECV状态,将客户端IP加入等待列表,
并重发SYN+ACK报文。
重发一般进行3-5次,大约间隔30秒左右轮询一次等待列表重试所有客户端。
另一方面,服务器在自己发出了SYN+ACK报文后,会预分配资源为即将建立的TCP连接储存信息做准备,
这个资源在等待重试期间一直保留。更为重要的是,服务器资源有限,可以维护的SYN_RECV状态超过极限后就不再接受新的SYN报文,
也就是拒绝新的TCP连接建立。

然而著名的SYNC Flood的DDos攻击正式利用上述的failover-recovery机制。
攻击者伪装大量的IP地址给服务器发送SYN报文,由于伪造的IP地址几乎不可能存在,也就几乎没有设备会给服务器返回任何应答了。
因此,服务器将会维持一个庞大的等待列表,不停地重试发送SYN+ACK报文,同时占用着大量的资源无法释放。
更为关键的是,被攻击服务器的SYN_RECV队列被恶意的数据包占满,不再接受新的SYN请求,合法用户无法完成三次握手建立起TCP连接。
也就是说,这个服务器被SYN Flood拒绝服务了。

SYN Flood攻击大量消耗服务器的CPU、内存资源,并占满SYN等待队列。相应的,我们修改内核参数即可有效缓解。主要参数如下:

1
2
3
net.ipv4.tcp_syncookies = 1
net.ipv4.tcp_max_syn_backlog = 8192
net.ipv4.tcp_synack_retries = 2

分别为启用SYN Cookie、设置SYN最大队列长度以及设置SYN+ACK最大重试次数。
SYNC Cookie主要是在服务端缓冲基于时间种子的SYN号,只有客户端发送的SYN+ACK与缓冲完全匹配才完成握手,否则直接丢弃。
tcp_max_syn_backlog则是增加等待队列的长度。

Apache storm中的ack机制

Apache storm是首个真正意义上的流式处理引擎,在spark/Flink出现之前,是实时计算领域的一枝独秀。

storm中是没有checkpoint机制的,但storm以大名鼎鼎的ack算法来保证at least once语义,(在Trident出现之前,storm是没有办法保证exactly once语义的)。ack需要spout节点保存每条数据,当所有的计算节点处理完毕,再发送给spout节点,因此与Chandy-Lamport算法相比,每条数据都需要保存和反复发送,而状态和数据回滚需要用户来保证。

实时大数据处理,数据源源不断的流入系统。无法在一个线程中串行的处理并确认一条或一批数据。
strom采用的是异步并行处理的模式(这里以JStorm的实现分析),
当excutor节点(executor节点是storm的进程,spout和bolt都是executor启动的task线程)收到消息时,首先将消息压入disruptor队列,disruptor的消费者从队列中获取数据,执行转发或者计算。

strom引入ack机制来确保数据不丢,但是对系统整体架构也带来了很大的影响,那么问题来了:

  • 消息量大,如何保存消息?
  • 消息可能流过多个节点,如何保证每个节点都正确处理?
  • spout节点重启,如何确保消息不丢失?
  • 消息堆积,如何确保集群的稳定?

ack机制是如何巧妙解决这写问题的呢?

1
2
A xor A = 0.
A xor B … xor B xor A = 0,

其中每一个操作数出现且仅出现两次。

strom的ack机制,巧妙的利用了两个相同的值异或为0的原理.

理解下整个大体节奏分为几部分:

  • 步骤1和2 spout把一条信息同时发送给了bolt1和bolt2,步骤3表示spout emit成功后去 acker bolt里注册本次根消息,ack值设定为本次发送的消息对应的64位id的异或运算值,上图对应的是T1^T2。

  • 步骤4表示bolt1收到T1后,单条tuple被拆成了三条消息T3T4T5发送给bolt3。步骤6 bolt1在ack()方法调用时会向acker bolt提交T1^T3^T4^T5的ack值。

  • 步骤5和7的bolt都没有产生新消息,所以ack()的时候分别向acker bolt提交了T2 和T3^T4^T5的ack值。综上所述,本次spout产生的tuple树对应的ack值经过的运算为 T1^T2^T1^T3^T4^T5^T2^T3^T4^T5按照异或运算的规则,ack值最终正好归零。

  • 步骤8为acker bolt发现根spout最终对应的的ack是0以后认为所有衍生出来的数据都已经处理成功,它会通知对应的spout,spout会调用相应的ack方法。

storm这个机制的实现方式保证了无论一个tuple树有多少个节点,一个根消息对应的追踪ack值所占用的空间大小是固定的,极大地节约了内存空间。

通过ack机制,spout发出的每一条消息,都可以确定是被成功或失败处理。但是,需要备份每条消息,来确认消息是否处理完成,如果消息流过的每个节点都备份数据,总数据量将翻几倍。spout作为消息流入到topology的起点,在这里备份数据既可以节省内存,又可以验证整条链路。此外,Ack机制还常用于限流作用: 为了避免spout发送数据太快,而bolt处理太慢,常常设置pending数,当spout有等于或超过pending数的tuple没有收到ack或fail响应时,跳过执行nextTuple, 从而限制spout发送数据。

strom逐条发送逐条处理逐条ack,这也是吞吐量不及spark和flink。

checkpoint

通俗来讲: 就是在分布式系统中,通过状态的checkpoint来确保数据的高可用。

checkpoint俗称检查点,是指定时将数据快照保存到持久化存储介质中,来提供数据的可靠性和与增量文件结合快速恢复数据。

Hadoop NameNode 的checkpoint

NameNode负责管理Hadoop的元数据(workspace信息、blockMap信息、network topology等)信息,是HDFS的心脏。
checkpoint机制是NameNode数据故障恢复的方案。

HDFS namenode 1.x

HDFS 2.x 引入了HA来解决NameNode的单点问题,社区也涌现了多种共享内存方案来保存editlog,而namenode的元数据的数据结构几乎没有变化。

name node workspace 内存结构

workspace信息常驻内存,并定时checkpoint成fsimage文件, 当HDFS-Client发起修改文件目录的请求时,直接修改内存中的数据, 并将修改记录写到editlog文件中。可以将name node的workspace的维护过程简单理解为分布式系统中消息处理的过程,

Chandy-Lamport算法

在实时流式处理中,简单的使用checkpoint没办法保证exactly once语义,主要是由于在某一个时刻:

  1. 消息还在处理(没有合并到状态中),source接收数据的偏移量不能准确的与状态做到数据一致性。
  2. 每个子任务处理进度也难以统一。

理想情况下,停止接收新数据并排干整个流处理系统,再做checkpoint,才能保证数据一致性和exactly once。停机显然是不可能的!Chandy-Lamport算法使用巧妙的方法,在不停止流处理的前提下拿到每个子任务的状态并checkpoint下来。

著名的一致性算法 Paxos 的作者Leslie Lamport与Chandy合作发表了算法论文: Distributed snapshots: determining global states of distributed systems, 在该论文中提出了分布式快照算法: Chandy-Lamport

A snapshot algorithm is used to create a consistent snapshot of the global state of a distributed system. Due to the lack of globally shared memory and a global clock, this isn’t trivially possible.

Chandy-Lamport算法用于在缺乏全局共享内存和全局时钟的分布式系统中创建一致性的全局分布式快照。而这个算法正是1978年提出的Time, Clocks and the Ordering of Events in a Distributed System的直接应用。在分布式系统中,为了确保数据在不同计算节点的有序性,引入barrier机制,当相同的barrier到达每一个计算节点时,认为全局节点处理结束。

Chandy-Lamport算法将全局的状态简化为有限个节点以及节点与节点之间的channel组成,也就是有向图。节点是进程,边是channel;分布式系统中,进程运行在不同的物理机器上,一个分布式的系统中的全局状态由进程的状态和channel中的message组成,这些都是分布式快照要保存的内容。

因为是有向图,一个节点的channel包含了input channel和output channel,流经channel的数据源源不断,假设channel是FIFO队列,保证不重复,那么只需要保存每个节点的局部状态和input message的偏移量,合并起来就是全局的分布式快照。

Flink中的Chandy-Lamport算法

Chandy-Lamport算法在flink中用于实现at least once语义。具体工作流程如下:

  1. 在checkpoint触发时刻,Job Manager会往所有Source的流中放入一个barrier(图中三角形)。barrier包含当前checkpoint的ID

  2. 当barrier经过一个subtask时,即表示当前这个subtask处于checkpoint触发的“时刻”,他就会立即将barrier法往下游,并执行checkpoint方法将当前的state存入backend storage。图中Source1和Source2就是完成了checkpoint动作。

  3. 如果一个subtask有多个上游节点,这个subtask就需要等待所有上游发来的barrier都接收到,才能表示这个subtask到达了checkpoint触发“时刻”。但所有节点的barrier不一定一起到达,这时候就会面临“是否要对齐barrier”的问题(Barrier Alignment)。如图中的Task1.1,他有2个上游节点,Source1和Source2。假设Source1的barrier先到,这时候Task1.1就有2个选择:

  • 是马上把这个barrier发往下游并等待Source2的barrier来了再做checkpoint
  • 还是把Source1这边后续的event全都cache起来,等Source2的barrier来了,在做checkpoint,完了再继续处理Source1和Source2的event,当前Source1这边需要先处理cache里的event。

WAL

WAL是一种常见的故障恢复方式,如NameNode的元数据、HBase WAL、kafka消息中间件、SQLite WAL等。

“In computer science, write-ahead logging (WAL) is a family of techniques for providing atomicity and durability (two of the ACID properties) in database systems.”——维基百科

HBase中的WAL

这里介绍一下HBase WAL(write ahead log)机制,Hbase的RegionServer在处理数据插入和删除的过程中用来记录操作内容的一种日志。在每次Put、Delete等一条记录时,首先将其数据写入到RegionServer对应的HLog文件中去。

客户端向RegionServer端提交数据的时候,会先写入WAL日志,只有当WAL日志写入成功的时候,客户端才会被告诉提交数据成功。如果写WAL失败会告知客户端提交失败,这其实就是数据落地的过程。

在一个RegionServer上的所有Region都共享一个HLog,一次数据的提交先写入WAL,写入成功后,再写入menstore之中。当menstore的值达到一定的时候,就会形成一个个StoreFile。

WAL记载了每一个RegionServer对应的HLog。RegionServer1或者RegionServer1上某一个regiong挂掉了,都会迁移到其它的机器上处理,重新操作,进行恢复。

当RegionServer意外终止的时候,Master会通过Zookeeper感知到,Master首先会处理遗留下来的HLog文件,将其中不同Region的Log数据进行拆分,分别放到相应的Region目录下,然后再将实效的Region重新分配,领取到这些Regio你的RegionMaster发现有历史的HLog需要处理,因此会Replay HLog的数据到Memstore之中,然后flush数据到StoreFiles,完成数据的恢复。

飞行日志+补偿机制,也是常用的方法,如基于消息的分布式事务是保证最终一致性的方式之一、Quartz中的恢复执行等。

[参考文献]

  1. 深入浅出DDoS攻击防御
  2. 《Storm源码分析》
  3. Flink Checkpoint
  4. 什么是WAL
  5. Write-Ahead Logging in SQLite

Kappa架构,是Jay Kreps在2014年提出的,其原文:《质疑Lambda架构》如下:

Nathan Marz写了一篇受欢迎的博客文章,描述了他称之为Lambda架构如何击败CAP定理)的想法。Lambda架构是一种在MapReduce和Storm或类似系统上构建流处理应用程序的方法。事实证明,这是一个出人意料的流行想法,有一个专门的网站即将初版的书。由于我一直在使用KafkaSamza在LinkedIn参与构建实时数据处理基础设施,经常被问及Lambda架构。我想描述我的想法和经历。

0.1 什么是Lambda架构,如何成为Lambda架构?

Lambda架构看起来像这样:

Lambda架构

Lambda架构的工作方法是捕获不可变的记录序列,同时输入批处理系统和流处理系统。需要分别实现批系统和流系统的计算逻辑,在查询时将批和流的计算结果组合成完整的结果。

Lambda架构有很多变体,我有意的简化一下。例如,你可以在Kafka、Storm和Hadoop的各种类似系统中交换数据,人们经常使用两个不同的数据库来存储输出表,一个优化为实时,另一个优化为批量更新。

Lambda架构的目标是围绕复杂的异步转换构建的应用程序,这些转换需要低延迟(比如,几秒到几小时)运行。一个很好的例子是新闻推荐系统,它需要抓取各种新闻源,处理并规范化所有输入,然后对其进行索引、排序和存储以供服务。

Lambda架构针对围绕复杂的异步转换构建的应用程序,这些转换需要以低延迟运行(例如,几秒钟到几个小时)。一个很好的例子是新闻推荐系统,该系统需要抓取各种新闻来源,处理和格式化所有输入,然后对其进行索引、排名和存储以供服务。

我在LinkedIn参与建立了许多实时数据系统和管道。其中有些就是这种风格,但经过深思熟虑,这并不是我最喜欢的方法。我认为有必要描述一下这个架构的优点和缺点,并给出一个我更喜欢的选择。

0.2 Lambda架构的优点

数据不可变

Lambda架构强调保持输入数据不变。将数据转换建模为从原始输入开始的一系列物化步骤是有很多优点的。这也是大型MapReduce工作流易于处理的原因之一,因为它使您能够独立调试每个阶段。我认为这一思想可以很好地应用到流处理领域。我在这里写过一些关于捕获和转换不可变数据流的想法。

支持重新计算

这个体系结构突出了重新计算的问题。重新计算是流处理的关键挑战之一,但经常被忽略。通过“重新计算”,再次处理输入数据以修正结果。这是一个完全明显但经常被忽视的要求。代码逻辑不可避免的修改。所以,如果您有从输入流中派生输出数据的代码,则每当代码更改时,都需要重新计算输出以查看更改的效果。

代码逻辑修改

随着业务的发展,需要增加新的字段或删掉不需要的字段,或者修复发现的bug。代码逻辑不可避免的修改。构建实时系统的人往往不考虑这一点,所以根本无法快速发展。因为Lambda架构方便的解决了重新计算的问题,是比较大的有点。

流处理也可以精确并强壮

Lambda架构的设计,认为流式计算本质上是近似处理且不强壮,比批处理更容易丢失数据。事实并非如此,虽然现在的流处理框架不如MapReduce成熟,但没有理由不能像批处理系统那样提供如此强大的语义保证。

击败CAP定理

Lambda架构可以权衡不同的数据系统的混合起来以某种方式“击败了CAP定理”。长话短说,虽然流处理中肯定有延迟/可用性权衡,但Lambda架构是一个异步处理的架构,因此正在计算中的结果并不能马上与输入的数据保持一致,这并不能满足CAP定理

0.3 Lambda架构的缺点

Lambda架构的问题在于,维护两套在复杂的分布式系统的代码逻辑,并且无法解决。

编写Storm和Hadoop等分布式框架的程序是很复杂的,并且代码嵌入到运行框架,导致实现Lamdba架构具有极高的复杂性。

为什么不能改进流处理系统以处理在其目标域中设置的完整问题?一个提出的修复方法是具有抽象实时和批处理框架的语言或框架。您使用此更高级别的框架编写代码,然后将其“按住”到封面下的流处理或MapReduce。summingbird是一个这样做的框架。这绝对会使事情变得更好,但我认为它不解决问题。

为什么不能改进流处理系统来处理其目标域中的完整问题集?解决这个问题的一个方法是在实时和批框架的基础上抽象出语言或框架。使用高级框架编写代码,然后它“向下编译”到的流处理或MapReduce。Summingbird是一个执行此操作的框架。这肯定会让事情变得更好一点,但它并不能解决问题。

即使避免了两次编码,运行和调试两个系统的操作成本也很高。任何新的抽象层只适配了两个系统支持的功能。更糟糕的是,对于这个新框架,将会无法兼容Hadoop如此强大的丰富的工具和语言生态系统(Hive、Pig、Crunch、Cascading、Oozie等)。

打个比方,跨数据库的ORM实现真正透明是臭名昭著的困难。而这还只是在类似的系统上和相近的标准接口语言上进行抽象。如果要在几乎不稳定的分布式系统之上抽象化完全不同的编程范式的问题要困难得多。

0.4 我们的尝试

在LinkedIn上进行了几轮这样的尝试。构建了各种混合Hadoop架构,甚至是一个特定于域的API,允许代码在实时或Hadoop上“透明”的运行。这些方法奏效了,但没有一种方法非常愉快或富有成效。让代码在两个不同的系统中完美同步是非常非常困难的。旨在隐藏底层框架的API被证明是最泄露的抽象,最终还是需要深入的了解实时层和Hadoop。并且还附加了新的要求,当您调试问题或试图解释性能时,需要了解API将如何转换为这些底层系统

我的建议是:如果对延迟不敏感,请使用MapReduce等批处理框架;如果对延迟敏感,请使用流处理框架;但除非您绝对必须,否则不要同时尝试同时进行这两种处理。

那么,为什么Lambda架构令人兴奋呢?我认为原因是人们越来越需要构建复杂、低延迟的处理系统。他们拥有的两个不能完全解决他们的问题:可扩展高延迟的批处理系统可以处理历史数据和低延迟流处理系统无法重新计算。通过将这两样东西粘在一起,它们实际上可以构建一个工作解决方案。

从这个意义上说,尽管Lambda架构可能会很痛苦,但解决了一个被普遍忽视的重要问题。但我不认为这是大数据的新范式或未来。它只是受限于现成工具的临时状态。我也认为有更好的选择。

0.5 Kappa架构

作为设计基础架构的人,我认为最明显的问题是:为什么不能改进流处理系统以处理其目标域中设置的完整问题?你为什么需要粘在另一个系统上?为什么不能进行实时处理并在代码更改时支持重新计算?流处理系统已经存在并行性的概念; 为什么不仅仅通过增加并行性和重放历史数据来处理重新计算,非常快,非常快?答案是您可以执行此操作,如果您今天正在建立这种类型的系统,我认为这实际上是一个合理的替代架构。

当我与人们讨论这个问题时,他们有时会告诉我,流处理不适合历史数据的高吞吐量处理。但我认为这是一种直觉,主要基于他们使用的系统的局限性,这些系统要么规模差,要么无法保存历史数据。这让他们感到,流处理系统本质上是计算一些短暂流的结果,然后扔掉所有底层数据的东西。但这是没有理由应该如此。流处理的基本抽象是数据流DAG,它们与传统数据仓库(火山)中的基本抽象完全相同,也是MapReduce继任Tez的基本抽象。流处理只是该数据流模型的推广,它向最终用户公开中间结果和持续输出的Checkpoint。

那么,我们如何直接从我们的流处理工作中进行后处理呢?我更喜欢的方法实际上是愚蠢的简单:

  1. 使用Kafka或其他系统,允许您保留您想要重新处理的数据的完整日志,并允许多个订阅者。例如,如果您想重新处理长达30天的数据,请将您在Kafka中的保留时间设置为30天。
  2. 当您想进行后处理时,请启动流处理作业的第二个实例,该实例从保留数据的开头开始处理,但将此输出数据定向到新的输出表。
  3. 当第二个作业赶上时,切换应用程序以从新表读取。
  4. 停止旧版本的作业,并删除旧的输出表。

这个架构看起来像这样:

卡帕

与Lambda架构不同,在这种方法中,您只在代码更改时进行重新计算,并且实际上也需要重新计算结果。当然,重新计算的工作只是同一代码的改进版本,在同一框架上运行,使用相同的输入数据。当然,您会想在重新处理任务中增加并行性,以便它很快就能完成。

也许我们可以称之为Kappa架构,尽管这个想法可能太简单了,不值得使用希腊字母。

当然,可以进一步优化这一点。在许多情况下,您可以将两个输出表合并。然而,我认为两者在短期内都有一些好处。这允许您只需有一个按钮将应用程序重定向到旧表,即可立即恢复到旧逻辑。在特别重要的情况下(例如,您的广告定位标准),您可以使用自动A/B测试或回溯算法进行切换,以确保新代码的任何错误修复或代码改进都不会意外地与之前的版本相比降级。

请注意,这并不意味着您的数据无法转到HDFS;它只是意味着您不会在那里运行后处理。Kafka与Hadoop集成良好,因此将任何Kafka Topic落库到HDFS中都很容易。Hadoop中中保存流处理作业的输出和中间流是非常有用的,可以用于Hive等工具中的分析或其他离线数据处理流的输入。

我们记录了这种方法的实施,以及使用Samza的后处理架构的其他变体。

0.6 一些背景

Kafka维护如下有序日志:

Kafka_log

Kafka的Topic是这些日志的集合:

partitioned_log

消费这些数据的流处理消费者只会维护一个“offset”,这是它在每个分区上处理的上一次记录的日志条目号。因此,更改消费者返回和重新处理数据的位置就像用不同的偏移量重新启动作业一样简单。为同一数据添加第二个消费者只是另一个指向日志中不同位置的读者。

Kafka支持复制和容错,在廉价的硬件上运行,并每台机器存储许多TB的数据是很轻松的。因此,保留大量数据是一件非常自然和经济的事情,不会损害性能。LinkedIn在线存储了超过兆字节的Kafka存储空间,许多应用程序正好为此目的很好地利用了这种长期保留模式。

廉价的消费者和保留大量数据的能力使添加第二个“重新计算”的作业只是启动代码的第二个实例,但从日志中的不同位置开始。

这个设计不是偶然的。我们构建Kafka的目的是将其用作流处理的基础,我们完全想到了这种重新计算的模型。感兴趣的话,可以在这里找到更多关于Kafka的信息。

然而,从根本上说,没有什么能将这个想法与Kafka联系起来。您可以替换任何支持长期保留有序数据的系统(例如HDFS或某种数据库)。事实上,许多人熟悉名为事件采购CQRS的类似模式。当然,分布式数据库人员会告诉你,这只是对实例化视图维护的轻微重塑,正如他们很乐意提醒你的那样,他们很久以前就知道了*,索尼*。

0.7 Lambda架构与Kappa架构对比

我知道使用Samza作为流处理系统这种方法效果很好,因为LinkedIn是这样做的。但我不知道为什么它不应该在Storm或其他流处理系统中同样有效。我对Storm不够熟悉,无法了解实用性,所以如果其他人已经在这样做,我很乐意听到。无论如何,我认为一般想法是相当独立的。

这两种方法之间的效率和资源权衡有点令人扫地。Lambda架构需要一直运行后处理和实时处理,而我提议的只是需要在您需要重新处理时运行作业的第二份副本。然而,我的提案要求在输出数据库中暂时拥有2倍的存储空间,并要求一个支持重装大批量写入的数据库。在这两种情况下,重新处理的额外负载可能会平均。如果您有很多这样的工作,它们不会同时进行重新处理,因此在包含数十个此类工作的共享集群上,您可能会为在任何给定时间积极重新处理的少数工作额外预算容量的一小部分。

真正的优势根本不在于效率,而在于允许人们在单个处理框架上开发、测试、调试和操作他们的系统。因此,在简单性很重要的情况下,请将这种方法视为Lambda架构的替代方案。

NameNode 高可用整体架构概述

在 Hadoop 1.0 时代,Hadoop 的两大核心组件 HDFS NameNode 和 JobTracker 都存在着单点问题,这其中以 NameNode 的单点问题尤为严重。因为 NameNode 保存了整个 HDFS 的元数据信息,一旦 NameNode 挂掉,整个 HDFS 就无法访问,同时 Hadoop 生态系统中依赖于 HDFS 的各个组件,包括 MapReduce、Hive、Pig 以及 HBase 等也都无法正常工作,并且重新启动 NameNode 和进行数据恢复的过程也会比较耗时。这些问题在给 Hadoop 的使用者带来困扰的同时,也极大地限制了 Hadoop 的使用场景,使得 Hadoop 在很长的时间内仅能用作离线存储和离线计算,无法应用到对可用性和数据一致性要求很高的在线应用场景中。

所幸的是,在 Hadoop2.0 中,HDFS NameNode 和 YARN ResourceManger(JobTracker 在 2.0 中已经被整合到 YARN ResourceManger 之中) 的单点问题都得到了解决,经过多个版本的迭代和发展,目前已经能用于生产环境。HDFS NameNode 和 YARN ResourceManger 的高可用 (High Availability,HA) 方案基本类似,两者也复用了部分代码,但是由于 HDFS NameNode 对于数据存储和数据一致性的要求比 YARN ResourceManger 高得多,所以 HDFS NameNode 的高可用实现更为复杂一些,本文从内部实现的角度对 HDFS NameNode 的高可用机制进行详细的分析。

HDFS NameNode 的高可用整体架构如图 1 所示 (图片来源于参考文献 [1]):

img

从上图中,我们可以看出 NameNode 的高可用架构主要分为下面几个部分:

Active NameNode 和 Standby NameNode:两台 NameNode 形成互备,一台处于 Active 状态,为主 NameNode,另外一台处于 Standby 状态,为备 NameNode,只有主 NameNode 才能对外提供读写服务。

主备切换控制器 ZKFailoverController:

ZKFailoverController 作为独立的进程运行,对 NameNode 的主备切换进行总体控制。ZKFailoverController 能及时检测到 NameNode 的健康状况,在主 NameNode 故障时借助 Zookeeper 实现自动的主备选举和切换,当然 NameNode 目前也支持不依赖于 Zookeeper 的手动主备切换。检测NameNode和主从切换 healthMonitor和ActiveStandbyElector

Zookeeper 集群:为主备切换控制器提供主备选举支持。

共享存储系统:

共享存储系统是实现 NameNode 的高可用最为关键的部分,共享存储系统保存了 NameNode 在运行过程中所产生的 HDFS 的元数据。主 NameNode 和备NameNode 通过共享存储系统实现元数据同步。在进行主备切换的时候,新的主 NameNode 在确认元数据完全同步之后才能继续对外提供服务。(会不会没有同步完,新的选举就开始了)

DataNode 节点:

除了通过共享存储系统共享 HDFS 的元数据信息之外,主 NameNode 和备 NameNode 还需要共享 HDFS 的数据块和 DataNode 之间的映射关系DataNode 会同时向主 NameNode 和备 NameNode 上报数据块的位置信息。 name node包含的元数据信息

下面开始分别介绍 NameNode 的主备切换实现和共享存储系统的实现,在文章的最后会结合笔者的实践介绍一下在 NameNode 的高可用运维中的一些注意事项。

NameNode 的主备切换实现

NameNode 主备切换主要由 ZKFailoverControllerHealthMonitorActiveStandbyElector 这 3 个组件来协同实现:

HealthMonitor负责监听NameNode的状态,而ActiveStandbyElector负责主备切换

ZKFailoverController 作为 NameNode 机器上一个独立的进程启动 (在 hdfs 启动脚本之中的进程名为 zkfc),启动的时候会创建 HealthMonitor 和 ActiveStandbyElector 这两个主要的内部组件,ZKFailoverController 在创建 HealthMonitor 和 ActiveStandbyElector 的同时,也会向 HealthMonitor 和 ActiveStandbyElector 注册相应的回调方法。 回调方法分别用于几个场景:

  • 强制fench NameNode

HealthMonitor 主要负责检测 NameNode 的健康状态,如果检测到 NameNode 的状态发生变化,会回调ZKFailoverController 的相应方法进行自动的主备选举。

ActiveStandbyElector 主要负责完成自动的主备选举,内部封装了 Zookeeper 的处理逻辑,一旦 Zookeeper 主备选举完成,会回调 ZKFailoverController 的相应方法来进行 NameNode 的主备状态切换。

NameNode 实现主备切换的流程如图 2 所示,有以下几步:

HAServiceProtocol RPC与Hadoop RPC的异同

  1. HealthMonitor 初始化完成之后会启动内部的线程来定时调用对应 NameNode 的 HAServiceProtocol RPC 接口的方法,对 NameNode 的健康状态进行检测。
  2. HealthMonitor 如果检测到 NameNode 的健康状态发生变化,会回调 ZKFailoverController 注册的相应方法进行处理
  3. 如果 ZKFailoverController 判断需要进行主备切换,会首先使用 ActiveStandbyElector 来进行自动的主备选举。
  4. ActiveStandbyElector 与 Zookeeper 进行交互完成自动的主备选举。
  5. ActiveStandbyElector 在主备选举完成后,会回调 ZKFailoverController 的相应方法来通知当前的 NameNode 成为主 NameNode 或备 NameNode。
  6. ZKFailoverController 调用对应 NameNode 的 HAServiceProtocol RPC 接口的方法将 NameNode 转换为 Active 状态或 Standby 状态。
图 2.NameNode 的主备切换流程

img

下面分别对 HealthMonitor、ActiveStandbyElector 和 ZKFailoverController 的实现细节进行分析:

HealthMonitor 实现分析

ZKFailoverController 在初始化的时候会创建 HealthMonitor,HealthMonitor 在内部会启动一个线程来循环调用 NameNode 的 HAServiceProtocol RPC 接口的方法来检测 NameNode 的状态,并将状态的变化通过回调的方式来通知 ZKFailoverController。

HealthMonitor 主要检测 NameNode 的两类状态,分别是 HealthMonitor.State 和 HAServiceStatus。HealthMonitor.State 是通过 HAServiceProtocol RPC 接口的 monitorHealth 方法来获取的,反映了 NameNode 节点的健康状况,主要是磁盘存储资源是否充足。HealthMonitor.State 包括下面几种状态:

  • •INITIALIZING:HealthMonitor 在初始化过程中,还没有开始进行健康状况检测;
  • •SERVICE_HEALTHY:NameNode 状态正常;
  • •SERVICE_NOT_RESPONDING:调用 NameNode 的 monitorHealth 方法调用无响应或响应超时;
  • •SERVICE_UNHEALTHY:NameNode 还在运行,但是 monitorHealth 方法返回状态不正常,磁盘存储资源不足;
  • •HEALTH_MONITOR_FAILED:HealthMonitor 自己在运行过程中发生了异常,不能继续检测 NameNode 的健康状况,会导致 ZKFailoverController 进程退出;

HealthMonitor.State 在状态检测之中起主要的作用,在 HealthMonitor.State 发生变化的时候,HealthMonitor 会回调 ZKFailoverController 的相应方法来进行处理,具体处理见后文 ZKFailoverController 部分所述。

而 HAServiceStatus 则是通过 HAServiceProtocol RPC 接口的 getServiceStatus 方法来获取的,主要反映的是 NameNode 的 HA 状态,包括:

  • •INITIALIZING:NameNode 在初始化过程中;
  • •ACTIVE:当前 NameNode 为主 NameNode;
  • •STANDBY:当前 NameNode 为备 NameNode;
  • •STOPPING:当前 NameNode 已停止;

HAServiceStatus 在状态检测之中只是起辅助的作用,在 HAServiceStatus 发生变化时,HealthMonitor 也会回调 ZKFailoverController 的相应方法来进行处理,具体处理见后文 ZKFailoverController 部分所述。

ActiveStandbyElector 实现分析

Namenode(包括 YARN ResourceManager) 的主备选举是通过 ActiveStandbyElector 来完成的,ActiveStandbyElector 主要是利用了 Zookeeper 的写一致性和临时节点机制,具体的主备选举实现如下:

创建锁节点

如果 HealthMonitor 检测到对应的 NameNode 的状态正常,那么表示这个 NameNode 有资格参加 Zookeeper 的主备选举。如果目前还没有进行过主备选举的话,那么相应的 ActiveStandbyElector 就会发起一次主备选举,尝试在 Zookeeper 上创建一个路径为/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 的临时节点 (${dfs.nameservices} 为 Hadoop 的配置参数 dfs.nameservices 的值,下同),Zookeeper 的写一致性会保证最终只会有一个 ActiveStandbyElector 创建成功,那么创建成功的 ActiveStandbyElector 对应的 NameNode 就会成为主 NameNode,ActiveStandbyElector 会回调 ZKFailoverController 的方法进一步将对应的 NameNode 切换为 Active 状态。而创建失败的 ActiveStandbyElector 对应的 NameNode 成为备 NameNode,ActiveStandbyElector 会回调 ZKFailoverController 的方法进一步将对应的 NameNode 切换为 Standby 状态。

注册 Watcher 监听

不管创建/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点是否成功,ActiveStandbyElector 随后都会向 Zookeeper 注册一个 Watcher 来监听这个节点的状态变化事件,ActiveStandbyElector 主要关注这个节点的 NodeDeleted 事件。

自动触发主备选举

如果 Active NameNode 对应的 HealthMonitor 检测到 NameNode 的状态异常时, ZKFailoverController 会主动删除当前在 Zookeeper 上建立的临时节点/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock, 这样处于 Standby 状态的 NameNode 的 ActiveStandbyElector 注册的监听器就会收到这个节点的 NodeDeleted 事件。收到这个事件之后,会马上再次进入到创建/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点的流程,如果创建成功,这个本来处于 Standby 状态的 NameNode 就选举为主 NameNode 并随后开始切换为 Active 状态。

当然,如果是 Active 状态的 NameNode 所在的机器整个宕掉的话,那么根据 Zookeeper 的临时节点特性,/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点会自动被删除,从而也会自动进行一次主备切换。

防止脑裂

Zookeeper 在工程实践的过程中经常会发生的一个现象就是 Zookeeper 客户端“假死”,所谓的“假死”是指如果 Zookeeper 客户端机器负载过高或者正在进行 JVM Full GC,那么可能会导致 Zookeeper 客户端到 Zookeeper 服务端的心跳不能正常发出,一旦这个时间持续较长,超过了配置的 Zookeeper Session Timeout 参数的话,Zookeeper 服务端就会认为客户端的 session 已经过期从而将客户端的 Session 关闭。“假死”有可能引起分布式系统常说的双主或脑裂 (brain-split) 现象。具体到本文所述的 NameNode,假设 NameNode1 当前为 Active 状态,NameNode2 当前为 Standby 状态。如果某一时刻 NameNode1 对应的 ZKFailoverController 进程发生了“假死”现象,那么 Zookeeper 服务端会认为 NameNode1 挂掉了,根据前面的主备切换逻辑,NameNode2 会替代 NameNode1 进入 Active 状态。但是此时 NameNode1 可能仍然处于 Active 状态正常运行,即使随后 NameNode1 对应的 ZKFailoverController 因为负载下降或者 Full GC 结束而恢复了正常,感知到自己和 Zookeeper 的 Session 已经关闭,但是由于网络的延迟以及 CPU 线程调度的不确定性,仍然有可能会在接下来的一段时间窗口内 NameNode1 认为自己还是处于 Active 状态。这样 NameNode1 和 NameNode2 都处于 Active 状态,都可以对外提供服务。这种情况对于 NameNode 这类对数据一致性要求非常高的系统来说是灾难性的,数据会发生错乱且无法恢复。Zookeeper 社区对这种问题的解决方法叫做 fencing,中文翻译为隔离,也就是想办法把旧的 Active NameNode 隔离起来,使它不能正常对外提供服务。

ActiveStandbyElector 为了实现 fencing,会在成功创建 Zookeeper 节点 hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 从而成为 Active NameNode 之后,创建另外一个路径为/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb 的持久节点,这个节点里面保存了这个 Active NameNode 的地址信息。Active NameNode 的 ActiveStandbyElector 在正常的状态下关闭 Zookeeper Session 的时候 (注意由于/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 是临时节点,也会随之删除),会一起删除节点/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb。但是如果 ActiveStandbyElector 在异常的状态下 Zookeeper Session 关闭 (比如前述的 Zookeeper 假死),那么由于/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb 是持久节点,会一直保留下来。后面当另一个 NameNode 选主成功之后,会注意到上一个 Active NameNode 遗留下来的这个节点,从而会回调 ZKFailoverController 的方法对旧的 Active NameNode 进行 fencing,具体处理见后文 ZKFailoverController 部分所述。

ZKFailoverController 实现分析

ZKFailoverController 在创建 HealthMonitor 和 ActiveStandbyElector 的同时,会向 HealthMonitor 和 ActiveStandbyElector 注册相应的回调函数,ZKFailoverController 的处理逻辑主要靠 HealthMonitor 和 ActiveStandbyElector 的回调函数来驱动。

对 HealthMonitor 状态变化的处理

如前所述,HealthMonitor 会检测 NameNode 的两类状态,HealthMonitor.State 在状态检测之中起主要的作用,ZKFailoverController 注册到 HealthMonitor 上的处理 HealthMonitor.State 状态变化的回调函数主要关注 SERVICE_HEALTHY、SERVICE_NOT_RESPONDING 和 SERVICE_UNHEALTHY 这 3 种状态:

  • •如果检测到状态为 SERVICE_HEALTHY,表示当前的 NameNode 有资格参加 Zookeeper 的主备选举,如果目前还没有进行过主备选举的话,ZKFailoverController 会调用 ActiveStandbyElector 的 joinElection 方法发起一次主备选举。
  • •如果检测到状态为 SERVICE_NOT_RESPONDING 或者是 SERVICE_UNHEALTHY,就表示当前的 NameNode 出现问题了,ZKFailoverController 会调用 ActiveStandbyElector 的 quitElection 方法删除当前已经在 Zookeeper 上建立的临时节点退出主备选举,这样其它的 NameNode 就有机会成为主 NameNode。

而 HAServiceStatus 在状态检测之中仅起辅助的作用,在 HAServiceStatus 发生变化时,ZKFailoverController 注册到 HealthMonitor 上的处理 HAServiceStatus 状态变化的回调函数会判断 NameNode 返回的 HAServiceStatus 和 ZKFailoverController 所期望的是否一致,如果不一致的话,ZKFailoverController 也会调用 ActiveStandbyElector 的 quitElection 方法删除当前已经在 Zookeeper 上建立的临时节点退出主备选举。

对 ActiveStandbyElector 主备选举状态变化的处理

在 ActiveStandbyElector 的主备选举状态发生变化时,会回调 ZKFailoverController 注册的回调函数来进行相应的处理:

  • •如果 ActiveStandbyElector 选主成功,那么 ActiveStandbyElector 对应的 NameNode 成为主 NameNode,ActiveStandbyElector 会回调 ZKFailoverController 的 becomeActive 方法,这个方法通过调用对应的 NameNode 的 HAServiceProtocol RPC 接口的 transitionToActive 方法,将 NameNode 转换为 Active 状态。
  • •如果 ActiveStandbyElector 选主失败,那么 ActiveStandbyElector 对应的 NameNode 成为备 NameNode,ActiveStandbyElector 会回调 ZKFailoverController 的 becomeStandby 方法,这个方法通过调用对应的 NameNode 的 HAServiceProtocol RPC 接口的 transitionToStandby 方法,将 NameNode 转换为 Standby 状态。
  • •如果 ActiveStandbyElector 选主成功之后,发现了上一个 Active NameNode 遗留下来的/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb 节点 (见“ActiveStandbyElector 实现分析”一节“防止脑裂”部分所述),那么 ActiveStandbyElector 会首先回调 ZKFailoverController 注册的 fenceOldActive 方法,尝试对旧的 Active NameNode 进行 fencing,在进行 fencing 的时候,会执行以下的操作:
  1. 1.首先尝试调用这个旧 Active NameNode 的 HAServiceProtocol RPC 接口的 transitionToStandby 方法,看能不能把它转换为 Standby 状态。
  2. 2.如果 transitionToStandby 方法调用失败,那么就执行 Hadoop 配置文件之中预定义的隔离措施,Hadoop 目前主要提供两种隔离措施,通常会选择 sshfence:
  • •sshfence:通过 SSH 登录到目标机器上,执行命令 fuser 将对应的进程杀死;
  • •shellfence:执行一个用户自定义的 shell 脚本来将对应的进程隔离;

只有在成功地执行完成 fencing 之后,选主成功的 ActiveStandbyElector 才会回调 ZKFailoverController 的 becomeActive 方法将对应的 NameNode 转换为 Active 状态,开始对外提供服务。

NameNode 的共享存储实现

过去几年中 Hadoop 社区涌现过很多的 NameNode 共享存储方案,比如 shared NAS+NFS、BookKeeper、BackupNode 和 QJM(Quorum Journal Manager) 等等。目前社区已经把由 Clouderea 公司实现的基于 QJM 的方案合并到 HDFS 的 trunk 之中并且作为默认的共享存储实现,本部分只针对基于 QJM 的共享存储方案的内部实现原理进行分析。为了理解 QJM 的设计和实现,首先要对 NameNode 的元数据存储结构有所了解。

NameNode 的元数据存储概述

一个典型的 NameNode 的元数据存储目录结构如图 3 所示 (图片来源于参考文献 [4]),这里主要关注其中的 EditLog 文件和 FSImage 文件:

图 3 .NameNode 的元数据存储目录结构

img

NameNode 在执行 HDFS 客户端提交的创建文件或者移动文件这样的写操作的时候,会首先把这些操作记录在 EditLog 文件之中,然后再更新内存中的文件系统镜像。内存中的文件系统镜像用于 NameNode 向客户端提供读服务,而 EditLog 仅仅只是在数据恢复的时候起作用。记录在 EditLog 之中的每一个操作又称为一个事务,每个事务有一个整数形式的事务 id 作为编号。EditLog 会被切割为很多段,每一段称为一个 Segment。正在写入的 EditLog Segment 处于 in-progress 状态,其文件名形如 edits_inprogress_${start_txid},其中${start_txid} 表示这个 segment 的起始事务 id,例如上图中的 edits_inprogress_0000000000000000020。而已经写入完成的 EditLog Segment 处于 finalized 状态,其文件名形如 edits_${start_txid}-${end_txid},其中${start_txid} 表示这个 segment 的起始事务 id,${end_txid} 表示这个 segment 的结束事务 id,例如上图中的 edits_0000000000000000001-0000000000000000019。

NameNode 会定期对内存中的文件系统镜像进行 checkpoint 操作,在磁盘上生成 FSImage 文件,FSImage 文件的文件名形如 fsimage_${end_txid},其中${end_txid} 表示这个 fsimage 文件的结束事务 id,例如上图中的 fsimage_0000000000000000020。在 NameNode 启动的时候会进行数据恢复,首先把 FSImage 文件加载到内存中形成文件系统镜像,然后再把 EditLog 之中 FsImage 的结束事务 id 之后的 EditLog 回放到这个文件系统镜像上。

基于 QJM 的共享存储系统的总体架构

基于 QJM 的共享存储系统主要用于保存 EditLog,并不保存 FSImage 文件。FSImage 文件还是在 NameNode 的本地磁盘上。QJM 共享存储的基本思想来自于 Paxos 算法 (参见参考文献 [3]),采用多个称为 JournalNode 的节点组成的 JournalNode 集群来存储 EditLog。每个 JournalNode 保存同样的 EditLog 副本。每次 NameNode 写 EditLog 的时候,除了向本地磁盘写入 EditLog 之外,也会并行地向 JournalNode 集群之中的每一个 JournalNode 发送写请求,只要大多数 (majority) 的 JournalNode 节点返回成功就认为向 JournalNode 集群写入 EditLog 成功。如果有 2N+1 台 JournalNode,那么根据大多数的原则,最多可以容忍有 N 台 JournalNode 节点挂掉。

基于 QJM 的共享存储系统的内部实现架构图如图 4 所示,主要包含下面几个主要的组件:

图 4 . 基于 QJM 的共享存储系统的内部实现架构图

img

FSEditLog:这个类封装了对 EditLog 的所有操作,是 NameNode 对 EditLog 的所有操作的入口。

JournalSet: 这个类封装了对本地磁盘和 JournalNode 集群上的 EditLog 的操作,内部包含了两类 JournalManager,一类为 FileJournalManager,用于实现对本地磁盘上 EditLog 的操作。一类为 QuorumJournalManager,用于实现对 JournalNode 集群上共享目录的 EditLog 的操作。FSEditLog 只会调用 JournalSet 的相关方法,而不会直接使用 FileJournalManager 和 QuorumJournalManager。

FileJournalManager:封装了对本地磁盘上的 EditLog 文件的操作,不仅 NameNode 在向本地磁盘上写入 EditLog 的时候使用 FileJournalManager,JournalNode 在向本地磁盘写入 EditLog 的时候也复用了 FileJournalManager 的代码和逻辑。

QuorumJournalManager:封装了对 JournalNode 集群上的 EditLog 的操作,它会根据 JournalNode 集群的 URI 创建负责与 JournalNode 集群通信的类 AsyncLoggerSet, QuorumJournalManager 通过 AsyncLoggerSet 来实现对 JournalNode 集群上的 EditLog 的写操作,对于读操作,QuorumJournalManager 则是通过 Http 接口从 JournalNode 上的 JournalNodeHttpServer 读取 EditLog 的数据。

AsyncLoggerSet:内部包含了与 JournalNode 集群进行通信的 AsyncLogger 列表,每一个 AsyncLogger 对应于一个 JournalNode 节点,另外 AsyncLoggerSet 也包含了用于等待大多数 JournalNode 返回结果的工具类方法给 QuorumJournalManager 使用。

AsyncLogger:具体的实现类是 IPCLoggerChannel,IPCLoggerChannel 在执行方法调用的时候,会把调用提交到一个单线程的线程池之中,由线程池线程来负责向对应的 JournalNode 的 JournalNodeRpcServer 发送 RPC 请求。

JournalNodeRpcServer:运行在 JournalNode 节点进程中的 RPC 服务,接收 NameNode 端的 AsyncLogger 的 RPC 请求。

JournalNodeHttpServer:运行在 JournalNode 节点进程中的 Http 服务,用于接收处于 Standby 状态的 NameNode 和其它 JournalNode 的同步 EditLog 文件流的请求。

下面对基于 QJM 的共享存储系统的两个关键性问题同步数据和恢复数据进行详细分析。

基于 QJM 的共享存储系统的数据同步机制分析

Active NameNode 和 StandbyNameNode 使用 JouranlNode 集群来进行数据同步的过程如图 5 所示,Active NameNode 首先把 EditLog 提交到 JournalNode 集群,然后 Standby NameNode 再从 JournalNode 集群定时同步 EditLog:

图 5 . 基于 QJM 的共享存储的数据同步机制

img

Active NameNode 提交 EditLog 到 JournalNode 集群

当处于 Active 状态的 NameNode 调用 FSEditLog 类的 logSync 方法来提交 EditLog 的时候,会通过 JournalSet 同时向本地磁盘目录和 JournalNode 集群上的共享存储目录写入 EditLog。写入 JournalNode 集群是通过并行调用每一个 JournalNode 的 QJournalProtocol RPC 接口的 journal 方法实现的,如果对大多数 JournalNode 的 journal 方法调用成功,那么就认为提交 EditLog 成功,否则 NameNode 就会认为这次提交 EditLog 失败。提交 EditLog 失败会导致 Active NameNode 关闭 JournalSet 之后退出进程,留待处于 Standby 状态的 NameNode 接管之后进行数据恢复。

从上面的叙述可以看出,Active NameNode 提交 EditLog 到 JournalNode 集群的过程实际上是同步阻塞的,但是并不需要所有的 JournalNode 都调用成功,只要大多数 JournalNode 调用成功就可以了。如果无法形成大多数,那么就认为提交 EditLog 失败,NameNode 停止服务退出进程。如果对应到分布式系统的 CAP 理论的话,虽然采用了 Paxos 的“大多数”思想对 C(consistency,一致性) 和 A(availability,可用性) 进行了折衷,但还是可以认为 NameNode 选择了 C 而放弃了 A,这也符合 NameNode 对数据一致性的要求。

Standby NameNode 从 JournalNode 集群同步 EditLog

当 NameNode 进入 Standby 状态之后,会启动一个 EditLogTailer 线程。这个线程会定期调用 EditLogTailer 类的 doTailEdits 方法从 JournalNode 集群上同步 EditLog,然后把同步的 EditLog 回放到内存之中的文件系统镜像上 (并不会同时把 EditLog 写入到本地磁盘上)。

这里需要关注的是:从 JournalNode 集群上同步的 EditLog 都是处于 finalized 状态的 EditLog Segment。“NameNode 的元数据存储概述”一节说过 EditLog Segment 实际上有两种状态,处于 in-progress 状态的 Edit Log 当前正在被写入,被认为是处于不稳定的中间态,有可能会在后续的过程之中发生修改,比如被截断。Active NameNode 在完成一个 EditLog Segment 的写入之后,就会向 JournalNode 集群发送 finalizeLogSegment RPC 请求,将完成写入的 EditLog Segment finalized,然后开始下一个新的 EditLog Segment。一旦 finalizeLogSegment 方法在大多数的 JournalNode 上调用成功,表明这个 EditLog Segment 已经在大多数的 JournalNode 上达成一致。一个 EditLog Segment 处于 finalized 状态之后,可以保证它再也不会变化。

从上面描述的过程可以看出,虽然 Active NameNode 向 JournalNode 集群提交 EditLog 是同步的,但 Standby NameNode 采用的是定时从 JournalNode 集群上同步 EditLog 的方式,那么 Standby NameNode 内存中文件系统镜像有很大的可能是落后于 Active NameNode 的,所以 Standby NameNode 在转换为 Active NameNode 的时候需要把落后的 EditLog 补上来。

基于 QJM 的共享存储系统的数据恢复机制分析

处于 Standby 状态的 NameNode 转换为 Active 状态的时候,有可能上一个 Active NameNode 发生了异常退出,那么 JournalNode 集群中各个 JournalNode 上的 EditLog 就可能会处于不一致的状态,所以首先要做的事情就是让 JournalNode 集群中各个节点上的 EditLog 恢复为一致。另外如前所述,当前处于 Standby 状态的 NameNode 的内存中的文件系统镜像有很大的可能是落后于旧的 Active NameNode 的,所以在 JournalNode 集群中各个节点上的 EditLog 达成一致之后,接下来要做的事情就是从 JournalNode 集群上补齐落后的 EditLog。只有在这两步完成之后,当前新的 Active NameNode 才能安全地对外提供服务。

补齐落后的 EditLog 的过程复用了前面描述的 Standby NameNode 从 JournalNode 集群同步 EditLog 的逻辑和代码,最终调用 EditLogTailer 类的 doTailEdits 方法来完成 EditLog 的补齐。使 JournalNode 集群上的 EditLog 达成一致的过程是一致性算法 Paxos 的典型应用场景,QJM 对这部分的处理可以看做是 Single Instance Paxos(参见参考文献 [3]) 算法的一个实现,在达成一致的过程中,Active NameNode 和 JournalNode 集群之间的交互流程如图 6 所示,具体描述如下:

图 6.Active NameNode 和 JournalNode 集群的交互流程图

img

生成一个新的 Epoch

Epoch 是一个单调递增的整数,用来标识每一次 Active NameNode 的生命周期,每发生一次 NameNode 的主备切换,Epoch 就会加 1。这实际上是一种 fencing 机制,为什么需要 fencing 已经在前面“ActiveStandbyElector 实现分析”一节的“防止脑裂”部分进行了说明。产生新 Epoch 的流程与 Zookeeper 的 ZAB(Zookeeper Atomic Broadcast) 协议在进行数据恢复之前产生新 Epoch 的过程完全类似:

    Active NameNode 首先向 JournalNode 集群发送 getJournalState RPC 请求,每个 JournalNode 会返回自己保存的最近的那个 Epoch(代码中叫 lastPromisedEpoch)。

    NameNode 收到大多数的 JournalNode 返回的 Epoch 之后,在其中选择最大的一个加 1 作为当前的新 Epoch,然后向各个 JournalNode 发送 newEpoch RPC 请求,把这个新的 Epoch 发给各个 JournalNode。

    每一个 JournalNode 在收到新的 Epoch 之后,首先检查这个新的 Epoch 是否比它本地保存的 lastPromisedEpoch 大,如果大的话就把 lastPromisedEpoch 更新为这个新的 Epoch,并且向 NameNode 返回它自己的本地磁盘上最新的一个 EditLogSegment 的起始事务 id,为后面的数据恢复过程做好准备。如果小于或等于的话就向 NameNode 返回错误。

    NameNode 收到大多数 JournalNode 对 newEpoch 的成功响应之后,就会认为生成新的 Epoch 成功。

在生成新的 Epoch 之后,每次 NameNode 在向 JournalNode 集群提交 EditLog 的时候,都会把这个 Epoch 作为参数传递过去。每个 JournalNode 会比较传过来的 Epoch 和它自己保存的 lastPromisedEpoch 的大小,如果传过来的 epoch 的值比它自己保存的 lastPromisedEpoch 小的话,那么这次写相关操作会被拒绝。一旦大多数 JournalNode 都拒绝了这次写操作,那么这次写操作就失败了。如果原来的 Active NameNode 恢复正常之后再向 JournalNode 写 EditLog,那么因为它的 Epoch 肯定比新生成的 Epoch 小,并且大多数的 JournalNode 都接受了这个新生成的 Epoch,所以拒绝写入的 JournalNode 数目至少是大多数,这样原来的 Active NameNode 写 EditLog 就肯定会失败,失败之后这个 NameNode 进程会直接退出,这样就实现了对原来的 Active NameNode 的隔离了。

选择需要数据恢复的 EditLog Segment 的 id

需要恢复的 Edit Log 只可能是各个 JournalNode 上的最后一个 Edit Log Segment,如前所述,JournalNode 在处理完 newEpoch RPC 请求之后,会向 NameNode 返回它自己的本地磁盘上最新的一个 EditLog Segment 的起始事务 id,这个起始事务 id 实际上也作为这个 EditLog Segment 的 id。NameNode 会在所有这些 id 之中选择一个最大的 id 作为要进行数据恢复的 EditLog Segment 的 id。

向 JournalNode 集群发送 prepareRecovery RPC 请求

NameNode 接下来向 JournalNode 集群发送 prepareRecovery RPC 请求,请求的参数就是选出的 EditLog Segment 的 id。JournalNode 收到请求后返回本地磁盘上这个 Segment 的起始事务 id、结束事务 id 和状态 (in-progress 或 finalized)。

这一步对应于 Paxos 算法的 Phase 1a 和 Phase 1b(参见参考文献 [3]) 两步。Paxos 算法的 Phase1 是 prepare 阶段,这也与方法名 prepareRecovery 相对应。并且这里以前面产生的新的 Epoch 作为 Paxos 算法中的提案编号 (proposal number)。只要大多数的 JournalNode 的 prepareRecovery RPC 调用成功返回,NameNode 就认为成功。

选择进行同步的基准数据源,向 JournalNode 集群发送 acceptRecovery RPC 请求 NameNode 根据 prepareRecovery 的返回结果,选择一个 JournalNode 上的 EditLog Segment 作为同步的基准数据源。选择基准数据源的原则大致是:在 in-progress 状态和 finalized 状态的 Segment 之间优先选择 finalized 状态的 Segment。如果都是 in-progress 状态的话,那么优先选择 Epoch 比较高的 Segment(也就是优先选择更新的),如果 Epoch 也一样,那么优先选择包含的事务数更多的 Segment。

在选定了同步的基准数据源之后,NameNode 向 JournalNode 集群发送 acceptRecovery RPC 请求,将选定的基准数据源作为参数。JournalNode 接收到 acceptRecovery RPC 请求之后,从基准数据源 JournalNode 的 JournalNodeHttpServer 上下载 EditLog Segment,将本地的 EditLog Segment 替换为下载的 EditLog Segment。

这一步对应于 Paxos 算法的 Phase 2a 和 Phase 2b(参见参考文献 [3]) 两步。Paxos 算法的 Phase2 是 accept 阶段,这也与方法名 acceptRecovery 相对应。只要大多数 JournalNode 的 acceptRecovery RPC 调用成功返回,NameNode 就认为成功。

向 JournalNode 集群发送 finalizeLogSegment RPC 请求,数据恢复完成

上一步执行完成之后,NameNode 确认大多数 JournalNode 上的 EditLog Segment 已经从基准数据源进行了同步。接下来,NameNode 向 JournalNode 集群发送 finalizeLogSegment RPC 请求,JournalNode 接收到请求之后,将对应的 EditLog Segment 从 in-progress 状态转换为 finalized 状态,实际上就是将文件名从 edits_inprogress_${startTxid} 重命名为 edits_${startTxid}-${endTxid},见“NameNode 的元数据存储概述”一节的描述。

只要大多数 JournalNode 的 finalizeLogSegment RPC 调用成功返回,NameNode 就认为成功。此时可以保证 JournalNode 集群的大多数节点上的 EditLog 已经处于一致的状态,这样 NameNode 才能安全地从 JournalNode 集群上补齐落后的 EditLog 数据。

需要注意的是,尽管基于 QJM 的共享存储方案看起来理论完备,设计精巧,但是仍然无法保证数据的绝对强一致,下面选取参考文献 [2] 中的一个例子来说明:

假设有 3 个 JournalNode:JN1、JN2 和 JN3,Active NameNode 发送了事务 id 为 151、152 和 153 的 3 个事务到 JournalNode 集群,这 3 个事务成功地写入了 JN2,但是在还没能写入 JN1 和 JN3 之前,Active NameNode 就宕机了。同时,JN3 在整个写入的过程中延迟较大,落后于 JN1 和 JN2。最终成功写入 JN1 的事务 id 为 150,成功写入 JN2 的事务 id 为 153,而写入到 JN3 的事务 id 仅为 125,如图 7 所示 (图片来源于参考文献 [2])。按照前面描述的只有成功地写入了大多数的 JournalNode 才认为写入成功的原则,显然事务 id 为 151、152 和 153 的这 3 个事务只能算作写入失败。在进行数据恢复的过程中,会发生下面两种情况:

图 7.JournalNode 集群写入的事务 id 情况

img

  • •如果随后的 Active NameNode 进行数据恢复时在 prepareRecovery 阶段收到了 JN2 的回复,那么肯定会以 JN2 对应的 EditLog Segment 为基准来进行数据恢复,这样最后在多数 JournalNode 上的 EditLog Segment 会恢复到事务 153。从恢复的结果来看,实际上可以认为前面宕机的 Active NameNode 对事务 id 为 151、152 和 153 的这 3 个事务的写入成功了。但是如果从 NameNode 自身的角度来看,这显然就发生了数据不一致的情况。
  • •如果随后的 Active NameNode 进行数据恢复时在 prepareRecovery 阶段没有收到 JN2 的回复,那么肯定会以 JN1 对应的 EditLog Segment 为基准来进行数据恢复,这样最后在多数 JournalNode 上的 EditLog Segment 会恢复到事务 150。在这种情况下,如果从 NameNode 自身的角度来看的话,数据就是一致的了。

事实上不光本文描述的基于 QJM 的共享存储方案无法保证数据的绝对一致,大家通常认为的一致性程度非常高的 Zookeeper 也会发生类似的情况,这也从侧面说明了要实现一个数据绝对一致的分布式存储系统的确非常困难。

NameNode 在进行状态转换时对共享存储的处理

下面对 NameNode 在进行状态转换的过程中对共享存储的处理进行描述,使得大家对基于 QJM 的共享存储方案有一个完整的了解,同时也作为本部分的总结。

NameNode 初始化启动,进入 Standby 状态

在 NameNode 以 HA 模式启动的时候,NameNode 会认为自己处于 Standby 模式,在 NameNode 的构造函数中会加载 FSImage 文件和 EditLog Segment 文件来恢复自己的内存文件系统镜像。在加载 EditLog Segment 的时候,调用 FSEditLog 类的 initSharedJournalsForRead 方法来创建只包含了在 JournalNode 集群上的共享目录的 JournalSet,也就是说,这个时候只会从 JournalNode 集群之中加载 EditLog,而不会加载本地磁盘上的 EditLog。另外值得注意的是,加载的 EditLog Segment 只是处于 finalized 状态的 EditLog Segment,而处于 in-progress 状态的 Segment 需要后续在切换为 Active 状态的时候,进行一次数据恢复过程,将 in-progress 状态的 Segment 转换为 finalized 状态的 Segment 之后再进行读取。

加载完 FSImage 文件和共享目录上的 EditLog Segment 文件之后,NameNode 会启动 EditLogTailer 线程和 StandbyCheckpointer 线程,正式进入 Standby 模式。如前所述,EditLogTailer 线程的作用是定时从 JournalNode 集群上同步 EditLog。而 StandbyCheckpointer 线程的作用其实是为了替代 Hadoop 1.x 版本之中的 Secondary NameNode 的功能,StandbyCheckpointer 线程会在 Standby NameNode 节点上定期进行 Checkpoint,将 Checkpoint 之后的 FSImage 文件上传到 Active NameNode 节点。

NameNode 从 Standby 状态切换为 Active 状态

当 NameNode 从 Standby 状态切换为 Active 状态的时候,首先需要做的就是停止它在 Standby 状态的时候启动的线程和相关的服务,包括上面提到的 EditLogTailer 线程和 StandbyCheckpointer 线程,然后关闭用于读取 JournalNode 集群的共享目录上的 EditLog 的 JournalSet,接下来会调用 FSEditLog 的 initJournalSetForWrite 方法重新打开 JournalSet。不同的是,这个 JournalSet 内部同时包含了本地磁盘目录和 JournalNode 集群上的共享目录。这些工作完成之后,就开始执行“基于 QJM 的共享存储系统的数据恢复机制分析”一节所描述的流程,调用 FSEditLog 类的 recoverUnclosedStreams 方法让 JournalNode 集群中各个节点上的 EditLog 达成一致。然后调用 EditLogTailer 类的 catchupDuringFailover 方法从 JournalNode 集群上补齐落后的 EditLog。最后打开一个新的 EditLog Segment 用于新写入数据,同时启动 Active NameNode 所需要的线程和服务。

NameNode 从 Active 状态切换为 Standby 状态

当 NameNode 从 Active 状态切换为 Standby 状态的时候,首先需要做的就是停止它在 Active 状态的时候启动的线程和服务,然后关闭用于读取本地磁盘目录和 JournalNode 集群上的共享目录的 EditLog 的 JournalSet。接下来会调用 FSEditLog 的 initSharedJournalsForRead 方法重新打开用于读取 JournalNode 集群上的共享目录的 JournalSet。这些工作完成之后,就会启动 EditLogTailer 线程和 StandbyCheckpointer 线程,EditLogTailer 线程会定时从 JournalNode 集群上同步 Edit Log。

NameNode 高可用运维中的注意事项

本节结合笔者的实践,从初始化部署和日常运维两个方面介绍一些在 NameNode 高可用运维中的注意事项。

初始化部署

如果在开始部署 Hadoop 集群的时候就启用 NameNode 的高可用的话,那么相对会比较容易。但是如果在采用传统的单 NameNode 的架构运行了一段时间之后,升级为 NameNode 的高可用架构的话,就要特别注意在升级的时候需要按照以下的步骤进行操作:

  1. 1.对 Zookeeper 进行初始化,创建 Zookeeper 上的/hadoop-ha/${dfs.nameservices} 节点。创建节点是为随后通过 Zookeeper 进行主备选举做好准备,在进行主备选举的时候会在这个节点下面创建子节点 (具体可参照“ActiveStandbyElector 实现分析”一节的叙述)。这一步通过在原有的 NameNode 上执行命令 hdfs zkfc -formatZK 来完成。
  2. 2.启动所有的 JournalNode,这通过脚本命令 hadoop-daemon.sh start journalnode 来完成。
  3. 3.对 JouranlNode 集群的共享存储目录进行格式化,并且将原有的 NameNode 本地磁盘上最近一次 checkpoint 操作生成 FSImage 文件 (具体可参照“NameNode 的元数据存储概述”一节的叙述) 之后的 EditLog 拷贝到 JournalNode 集群上的共享目录之中,这通过在原有的 NameNode 上执行命令 hdfs namenode -initializeSharedEdits 来完成。
  4. 4.启动原有的 NameNode 节点,这通过脚本命令 hadoop-daemon.sh start namenode 完成。
  5. 5.对新增的 NameNode 节点进行初始化,将原有的 NameNode 本地磁盘上最近一次 checkpoint 操作生成 FSImage 文件拷贝到这个新增的 NameNode 的本地磁盘上,同时需要验证 JournalNode 集群的共享存储目录上已经具有了这个 FSImage 文件之后的 EditLog(已经在第 3 步完成了)。这一步通过在新增的 NameNode 上执行命令 hdfs namenode -bootstrapStandby 来完成。
  6. 6.启动新增的 NameNode 节点,这通过脚本命令 hadoop-daemon.sh start namenode 完成。
  7. 7.在这两个 NameNode 上启动 zkfc(ZKFailoverController) 进程,谁通过 Zookeeper 选主成功,谁就是主 NameNode,另一个为备 NameNode。这通过脚本命令 hadoop-daemon.sh start zkfc 完成。

日常维护

笔者在日常的维护之中主要遇到过下面两种问题:

Zookeeper 过于敏感:Hadoop 的配置项中 Zookeeper 的 session timeout 的配置参数 ha.zookeeper.session-timeout.ms 的默认值为 5000,也就是 5s,这个值比较小,会导致 Zookeeper 比较敏感,可以把这个值尽量设置得大一些,避免因为网络抖动等原因引起 NameNode 进行无谓的主备切换。

单台 JouranlNode 故障时会导致主备无法切换:在理论上,如果有 3 台或者更多的 JournalNode,那么挂掉一台 JouranlNode 应该仍然可以进行正常的主备切换。但是笔者在某次 NameNode 重启的时候,正好赶上一台 JournalNode 挂掉宕机了,这个时候虽然某一台 NameNode 通过 Zookeeper 选主成功,但是这台被选为主的 NameNode 无法成功地从 Standby 状态切换为 Active 状态。事后追查原因发现,被选为主的 NameNode 卡在退出 Standby 状态的最后一步,这个时候它需要等待到 JournalNode 的请求全部完成之后才能退出。但是由于有一台 JouranlNode 宕机,到这台 JournalNode 的请求都积压在一起并且在不断地进行重试,同时在 Hadoop 的配置项中重试次数的默认值非常大,所以就会导致被选为主的 NameNode 无法及时退出 Standby 状态。这个问题主要是 Hadoop 内部的 RPC 通信框架的设计缺陷引起的,Hadoop HA 的源代码 IPCLoggerChannel 类中有关于这个问题的 TODO,但是截止到社区发布的 2.7.1 版本这个问题仍然存在。

相关主题

yarn作为目前最流行的分布式计算资源管理平台,为Hadoop MR、spark、Flink等提供了资源容器

Docker File

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
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
# 基础镜像
FROM hub.c.163.com/public/centos
RUN yum clean all
RUN yum install -y yum-plugin-ovl || true
# 安装基础的工具包
RUN yum install -y vim tar wget curl rsync bzip2 iptables tcpdump less telnet net-tools lsof sysstat cronie python-setuptools
RUN yum clean all

RUN cp -f /usr/share/zoneinfo/Asia/Shanghai /etc/localtime
EXPOSE 22
RUN mkdir -p /etc/supervisor/conf.d/
RUN echo [include] >> /etc/supervisord.conf
RUN echo 'files = /etc/supervisor/conf.d/*.conf' >> /etc/supervisord.conf
#COPY sshd.conf /etc/supervisor/conf.d/sshd.conf
CMD ["/usr/bin/supervisord"]



# 镜像的作者
MAINTAINER averyzhang
# 安装openssh-server和sudo软件包,并且将sshd的UsePAM参数设置成no
RUN yum install -y openssh-server sudo
RUN sed -i 's/UsePAM yes/UsePAM no/g' /etc/ssh/sshd_config
#安装openssh-clients
RUN yum install -y openssh-clients

# 添加测试用户root,密码abc.123,并且将此用户添加到sudoers里
RUN echo "root:abc.123" | chpasswd
RUN echo "root ALL=(ALL) ALL" >> /etc/sudoers
# 下面这两句比较特殊,
# 在centos6上必须要有,否则创建出来的容器sshd不能登录

# 为了避免文件已存在报错,首先删掉私钥文件
RUN rm -rf /etc/ssh/ssh_host_rsa_key
RUN rm -rf /etc/ssh/ssh_host_dsa_key

RUN ssh-keygen -t dsa -f /etc/ssh/ssh_host_dsa_key
RUN ssh-keygen -t rsa -f /etc/ssh/ssh_host_rsa_key

# 启动sshd服务并且暴露22端口
# RUN mkdir /var/run/sshd
EXPOSE 22
CMD ["/usr/sbin/sshd", "-D"]
###### 以上建立centos-ssh

# 安装java
RUN yum install -y java-1.8.0-openjdk.x86_64 java-1.8.0-openjdk-devel.x86_64
## 添加环境变量
ENV JAVA_HOME /usr/lib/jvm/jre-openjdk/
ENV JRE_HOME ${JAVA_HOME}
ENV PATH $JAVA_HOME/bin:$PATH
ENV CLASSPATH $CLASSPATH:.:$JAVA_HOME/lib

RUN echo "export JAVA_HOME=/usr/lib/jvm/jre-openjdk/" >> /etc/profile.d/java.sh
RUN echo "export JRE_HOME=${JAVA_HOME}" >> /etc/profile.d/java.sh
RUN echo "export PATH=$JAVA_HOME/bin:$PATH" >> /etc/profile.d/java.sh
RUN echo "export CLASSPATH=$CLASSPATH:.:$JAVA_HOME/lib" >> /etc/profile.d/java.sh


## 安装scala
RUN mkdir -p /data/
RUN wget https://downloads.lightbend.com/scala/2.10.7/scala-2.10.7.rpm
RUN yum install -y scala-2.10.7.rpm
ENV SCALA_HOME /usr/share/java
RUN echo "export SCALA_HOME=/usr/share/java" >> /etc/profile.d/java.sh


## 安装Hadoop
## RUN wget http://mirror.bit.edu.cn/apache/hadoop/common/hadoop-2.7.7/hadoop-2.7.7.tar.gz
ADD hadoop-2.7.7.tar /data/hadoop
# RUN tar zxvf /data/hadoop-2.7.7.tar -C /data/hadoop
# RUN ln -s /data/hadoop-2.7.7 /data/hadoop
ENV HADOOP_HOME /data/hadoop
ENV HADOOP_CONF_DIR ${HADOOP_HOME}/etc/hadoop
ENV YARN_HOME ${HADOOP_HOME}
ENV YARN_CONF_DIR ${YARN_HOME}/etc/hadoop



RUN echo "export HADOOP_HOME=/data/hadoop" >> /etc/profile.d/hadoop.sh
RUN echo "export HADOOP_CONF_DIR=${HADOOP_HOME}/etc/hadoop" >> /etc/profile.d/hadoop.sh
RUN echo "export YARN_HOME=${HADOOP_HOME}" >> /etc/profile.d/hadoop.sh
RUN echo "export YARN_CONF_DIR=${YARN_HOME}/etc/hadoop" >> /etc/profile.d/hadoop.sh

制作image

1
docker build -t="avery/centos-yarn" .

创建yarn节点

1
2
3
4
docker run --name yarn0 --hostname yarn0 -d -P -p 50070:50070 -p 8088:8088 avery/centos-yarn
docker run --name yarn1 --hostname yarn1 -d -P avery/centos-yarn
docker run --name yarn2 --hostname yarn2 -d -P avery/centos-yarn
docker run --name yarn3 --hostname yarn3 -d -P avery/centos-yarn

查看三个container的ip:

1
2
3
4
5
6
7
8
% docker inspect --format='{{.NetworkSettings.IPAddress}}' yarn0
172.17.0.2
% docker inspect --format='{{.NetworkSettings.IPAddress}}' yarn1
172.17.0.3
% docker inspect --format='{{.NetworkSettings.IPAddress}}' yarn2
172.17.0.4
% docker inspect --format='{{.NetworkSettings.IPAddress}}' yarn3
172.17.0.5
节点 备注 ip
yarn0 master 172.17.0.2
yarn1 slaver 172.17.0.3
yarn2 slaver 172.17.0.4
yarn3 slaver 172.17.0.5

查看端口

1
2
3
4
5
6
docker container ls -a                                        
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
e150f9141cbd avery/centos-yarn "/usr/sbin/sshd -D" 2 minutes ago Up 2 minutes 0.0.0.0:32771->22/tcp yarn3
9fe0eb86b3e9 avery/centos-yarn "/usr/sbin/sshd -D" 2 minutes ago Up 2 minutes 0.0.0.0:32770->22/tcp yarn2
1d5c0ea3d3b3 avery/centos-yarn "/usr/sbin/sshd -D" 2 minutes ago Up 2 minutes 0.0.0.0:32769->22/tcp yarn1
60a9cb47ff0c avery/centos-yarn "/usr/sbin/sshd -D" 2 minutes ago Up 2 minutes 0.0.0.0:8088->8088/tcp, 0.0.0.0:50070->50070/tcp, 0.0.0.0:32768->22/tcp yarn0

连接container

验证ssh连接

1
2
% ssh root@localhost -p 32774
# 输入密码 abc.123

使用exec

1
docker exec -it yarn0 /bin/bash

修改container的主机名

分别修改三个container的hosts

1
vi /etc/hosts 

添加下面配置

1
2
3
4
172.17.0.2      yarn0
172.17.0.3 yarn1
172.17.0.4 yarn2
172.17.0.5 yarn3

ssh 免密

设置ssh免密码登录
在yarn0上执行下面操作

1
2
3
4
5
6
7
8
9
10
cd  ~
mkdir .ssh
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i yarn0
ssh-copy-id -i yarn1
ssh-copy-id -i yarn2
ssh-copy-id -i yarn3

在yarn1上执行下面操作

1
2
3
4
5
6
7
8
cd  ~
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i yarn1
ssh-copy-id -i yarn2
ssh-copy-id -i yarn3

在yarn2上执行下面操作

1
2
3
4
5
6
7
8
9
cd  ~
cd .ssh
ssh-keygen -t rsa
# (一直按回车即可)
ssh-copy-id -i localhost
ssh-copy-id -i yarn0
ssh-copy-id -i yarn1
ssh-copy-id -i yarn2
ssh-copy-id -i yarn3

至此,Docker搭建Hadoop集群的准备工作

Hadoop 安装

登录yarn0,修改Hadoop配置

配置 Hadoop,cd ~/hadoop-2.7.2/etc/hadoop进入hadoop配置目录,需要配置有以下7个文件:hadoop-env.sh,yarn-env.sh,slaves,core-site.xml,hdfs-site.xml,maprd-site.xml,yarn-site.xml。

在hadoop-env.sh中配置JAVA_HOME

1
2
# The java implementation to use.
export JAVA_HOME=/usr/lib/jvm/jdk1.8.0_77

在yarn-env.sh中配置JAVA_HOME

1
2
# some Java parameters
export JAVA_HOME=/usr/lib/jvm/jdk1.8.0_77

在slaves中配置slave节点的ip或者host,

yarn1
yarn2
yarn3

修改core-site.xml

1
2
3
4
5
6
7
8
9
10
<configuration>
<property>
<name>fs.defaultFS</name>
<value>hdfs://yarn0:9000/</value>
</property>
<property>
<name>hadoop.tmp.dir</name>
<value>file:/home/fang//hadoop-2.7.2/tmp</value>
</property>
</configuration>

修改hdfs-site.xml

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
<configuration>
<property>
<name>dfs.namenode.secondary.http-address</name>
<value>yarn0:9001</value>
</property>
<property>
<name>dfs.namenode.name.dir</name>
<value>file:/home/fang/hadoop-2.7.2/dfs/name</value>
</property>
<property>
<name>dfs.datanode.data.dir</name>
<value>file:/home/fang/hadoop-2.7.2/dfs/data</value>
</property>
<property>
<name>dfs.replication</name>
<value>3</value>
</property>
</configuration>

修改mapred-site.xml

1
2
3
4
5
6
<configuration>
<property>
<name>mapreduce.framework.name</name>
<value>yarn</value>
</property>
</configuration>

修改yarn-site.xml

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
<configuration>
<property>
<name>yarn.nodemanager.aux-services</name>
<value>mapreduce_shuffle</value>
</property>
<property>
<name>yarn.nodemanager.aux-services.mapreduce.shuffle.class</name>
<value>org.apache.hadoop.mapred.ShuffleHandler</value>
</property>
<property>
<name>yarn.resourcemanager.address</name>
<value>yarn0:8032</value>
</property>
<property>
<name>yarn.resourcemanager.scheduler.address</name>
<value>yarn0:8030</value>
</property>
<property>
<name>yarn.resourcemanager.resource-tracker.address</name>
<value>yarn0:8035</value>
</property>
<property>
<name>yarn.resourcemanager.admin.address</name>
<value>yarn0:8033</value>
</property>
<property>
<name>yarn.resourcemanager.webapp.address</name>
<value>yarn0:8088</value>
</property>
</configuration>

将配置好的hadoop-2.7.2文件夹分发给所有slaves节点

1
2
3
scp -r hadoop/* root@yarn1:/data/hadoop
scp -r hadoop/* root@yarn2:/data/hadoop
scp -r hadoop/* root@yarn3:/data/hadoop

启动HDFS和yarn

格式化HDFS

1
2
3
4
bin/hadoop namenode -format    #格式化namenode
注:若格式化之后重新修改了配置文件,重新格式化之前需要删除tmp,dfs,logs文件夹。
sbin/start-dfs.sh              #启动dfs 
sbin/start-yarn.sh              #启动yarn

检查

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

[root@yarn0 hadoop]# jps
709 ResourceManager
984 Jps
346 NameNode
543 SecondaryNameNode
[root@yarn0 hadoop]# ssh yarn1
Last login: Sun Jun 2 23:21:51 2019 from yarn0
[root@yarn1 ~]# jps
305 NodeManager
449 Jps
200 DataNode
[root@yarn1 ~]# ssh yarn2
Last login: Sun Jun 2 23:22:18 2019 from yarn0
[root@yarn2 ~]# jps
193 DataNode
442 Jps
298 NodeManager
[root@yarn2 ~]# ssh yarn3
Last login: Sun Jun 2 23:22:25 2019 from yarn0
[root@yarn3 ~]# jps
178 DataNode
283 NodeManager
427 Jps

yarn web

浏览器打开 http://localhost:8088

搞定收工

1 Apache Parquet

1.1 概念

大规模分析型数据处理在互联网乃至其他行业中应用都已越来越广泛,尤其是当前已经可以用廉价的存储来收集、保存海量的业务数据情况下。如何让分析师和工程师便捷的利用这些数据也变得越来越重要。列式存储(Column-oriented Storage)是大数据场景面向分析型数据的主流存储方式。与行式存储相比,列存由于可以只提取部分数据列、同列同质数据拥有更好的编码及压缩方式,因此在 OLAP 场景下能提供更好的 IO 性能。

Apache Parquet 是由 Twitter 和 Cloudera 最先发起并合作开发的列存项目,也是 2010 年 Google 发表的 Dremel 论文中描述的内部列存格式的开源实现。和一些传统的列式存储(C-Store、MonetDB 等)系统相比,Dremel/Parquet 最大的贡献是支持嵌套格式数据(Nested Data)的列式存储。嵌套格式可以很自然的描述互联网和科学计算等领域的数据,Dremel/Parquet “原生”的支持嵌套格式数据减少了规则化、重新组合这些大规模数据的代价。

Parquet 的设计与计算框架、数据模型以及编程语言无关,可以与任意项目集成,因此应用广泛。目前已经是 Hadoop 大数据生态圈列式存储的事实标准。

1.2 原理

1.2.1 行存 VS 列存

例如,下图是拥有 A/B/C 3 个字段的简单示意表:

图片

在面向行的存储中,每列的数据依次排成一行,如下所示:

图片

而在面向列的存储中,相同列的数据存储在一起:

图片

显而易见,行存适用于数据整行读取场景,而列存适用于只读取部分列数据(如统计分析等)场景。

1.2.2 数据模型

1.2.2.1 (1).schema 协议

想要深入的了解 Parquet 存储格式首先需要理解它的数据模型。
Parquet 采用了一个类似 Google Protobuf 的协议来描述存储数据的 schema。
下面是 Parquet 数据 schema 的一个简单示例:

1
message AddressBook {

schema 的最上层是 message,里面可以包含一系列字段。
每个字段都拥有3个属性:重复性(repetition)、类型(type)以及名称(name)。
字段类型可以是一个 group 或者原子类型(如 int/boolean/string 等),group 可以用来表示数据的嵌套结构。
字段的重复性有三种情况:

  • required:有且只有一次
  • optional:0或1次
  • repeated:0或多次

这个模型非常的简洁。一些复杂的数据类型如:Map/List/Set 也可以用重复的字段(repeated fields) + groups 来表达,因此也就不用再单独定义这些类型。

采用 repeated field 表达 List 或者 Set 的示例:

图片

采用 repeated group(包含 key 和 value,其中 key 是 required) 来表达 Map 的示例:

图片

1.2.2.2 (2).列式存储格式

为了使数据能够按列存储,对于一条记录(Record),首先要将其按列(Column)进行拆分。对于扁平(Flat)结构数据,拆分比较直观,一个字段即对应一列,而嵌套格式数据会复杂些。
Dremel/Parquet 中,提出以树状层级的形式组织 schema 中的字段(Field),树的叶子结点对应一个原子类型字段,这样这个模型能同时覆盖扁平结构和嵌套结构数据(扁平结构只是嵌套结构的一种特例)。嵌套字段的完整路径使用简单的点分符号表示,如 contacts.name

AddressBook 例子以树状结构展示的样式:

图片

列存连续的存储一个字段的值,以便进行高效的编码压缩及快速的读取。Dremel 中行存 vs 列存的图示:

图片

1.2.2.3 (3).Repetition and Defination Levels

对于嵌套格式列存,除了按列拆分进行连续的存储,还需要能够“无损”的保留嵌套格式的结构化信息,以便正确的重建记录。

只有字段值不能表达清楚记录的结构。给定一个重复字段的两个值,我们不知道此值是在什么“级别”被重复的(比如,这些值是来自两个不同的记录,还是相同的记录中两个重复的值)。同样的,给出一个缺失的可选字段,我们不知道整个路径有多少字段被显示定义了。

Dremel 提出了 Repetition Level(重复级别)和 Definition Level(定义级别)两个概念,用以解决这个问题。并实现了记录中任意一个字段的恢复都不需要依赖其它字段,且可以对任意字段子集按原始嵌套格式进行重建。

图片

Repetition levels:用以表示在该字段路径上哪个节点进行了重复(at what repeated field in the field’s path the value has repeated)。

一个重复字段存储的列值,有可能来自不同记录,也可能由同一记录的不同层级节点重复导致。如上图中的 Code 字段,他在 r1 记录中出现了 3 次,分别是字段 Name 和 Language 重复导致的,其中 Language 先重复了 2 次,Name 字段再重复了 1 次。

Repetition Levels 采用数字代表重复节点的层级。根据树形层次结构,根结点为 0、下一层级为 1… 依次类推。根结点的重复暗含了记录的重复,也即 r=0 代表新记录的开始。required 和 optional 字段不需要 repetition level,只有可重复的字段需要。因此,上述 Code 字段的 repetition levels 范围为 0-2。当我们从上往下扫描 r1 记录,首先遇到 Code 的值是“en-us”,由于它之前没有该字段路径相关的字段出现,因此 r=0;其次遇到“en”,是 Language 层级重复导致的,r=2;最后遇到“en-gb”,是 Name 层级重复导致的,因此 r=1。所以,Code 字段的 repetition levels 在 r1 记录中是“0,2,1”。

需要注意的是,r1 记录中的第二个重复 Name,由于其不包含 Code 字段,为了区分“en-gb”值是来自记录中的第三个 Name 而不是第二个,我们需要在“en”和“en-gb”之间插入一个值“null”。由于它是 Name 级重复的,因此它的 r=1。另外还需要注意一些隐含信息,比如 Code 是 required 字段类型,因此一旦 Code 出现未定义,则隐含表明其上级 Language 也肯定未定义。

Definition Levels:用以表示该字段路径上有多少可选的字段实际进行了定义(how many fields in p that could be undefined (because they are optional or repeated) are actually present)。

光有 Repetition Levels 尚无法完全保留嵌套结构信息,考虑上述图中 r1 记录的 Backward 字段。由于 r1 中未定义 Backward 字段,因此我们插入一个“null”并设置 r=0。但 Backward 的上级 Links 字段在 r1 中显式的进行了定义,null 和 r=0 无法再表达出这一层信息。因此需要额外再添加 Definition Levels 定义记录可选字段出现的个数,Backward 的路径上出现 1 个可选字段 Links,因此它的 d=1。

有了 Definition Levels 我们就可以清楚的知道该值出现在字段路径的第几层,对未定义字段的 null 和字段实际的值为 null 也能进行区分。只有 optional 和 repeated 字段需要 Definition Levels 定义,因为 required 字段已经隐含了字段肯定被定义(这可以减少 Definition Levels 需要描述的数字,并在一定程度上节省后续的存储空间)。另外一些其他的隐含信息:如果 Definition Levels 小于路径中 optional + repeated 字段的数量,则该字段的值肯定为 null;Definition Levels 的值为 0 隐含了 Repeated Levels 也为 0(路径中没有 optional/repeated 字段或整个路径未定义)。

1.2.2.4 (4). striping and assembly 算法

现在把 Repetition LevelsDefinition Levels 两个概念一起考虑。还是沿用上述 AddressBook 例子。
下表显示了 AddressBook 中每个字段的最大重复和定义级别,并解释了为什么它们小于列的深度:

图片

假设这是两条真实的 AddressBook 数据:

1
AddressBook {

我们采用 contacts.phoneNumber 字段来演示一下拆解和重组记录的 striping and assembly 算法。

仅针对 contacts.phoneNumber 字段投影后,数据具有如下结构:

1
AddressBook {

计算可得该字段对应的数据如下(R =重复级别,D =定义级别):

图片

因此我们最终存储的记录数据如下:

1
contacts.phoneNumber: “555 987 6543”

使用图表展示(注意其中的 null 值并不会实际存储,原因如上所说只要 Definition Levels 小于其 max 值即隐含该字段值为 null):

图片

在重组该记录时,我们重复读取该字段的值:

1
R=0, D=2, Value = “555 987 6543”:

1.3 工程实现

Parquet 工程具体的实现。

1.3.1 Parquet 文件存储格式中的术语

  • Block (hdfs block):即指 HDFS Block,Parquet 的设计与 HDFS 完全兼容。Block 是 HDFS 文件存储的基本单位,HDFS 会维护一个 Block 的多个副本。在 Hadoop 1.x 版本中 Block 默认大小 64M,Hadoop 2.x 版本中默认大小为 128M。
  • File:HDFS 文件,保存了该文件的元数据信息,但可以不包含实际数据(由 Block 保存)。
  • Row group:按照行将数据划分为多个逻辑水平分区。一个 Row group(行组)由每个列的一个列块(Column Chunk)组成。
  • Column chunk:一个列的列块,分布在行组当中,并在文件中保证是连续的。
  • Page:一个列块切分成多个 Pages(页面),概念上讲,页面是 Parquet 中最小的基础单元(就压缩和编码方面而言)。一个列块中可以有多个类型的页面。

1.3.2 并行化执行的基本单元

  • MapReduce - File/Row Group(一个任务对应一个文件或一个行组)
  • IO - Column chunk(任务中的 IO 以列块为单位进行读取)
  • Encoding/Compression - Page(编码格式和压缩一次以一个页面为单位进行)

1.3.3 Parquet 文件格式

Parquet 文件格式是自解析的,采用 thrift 格式定义的文件 schema 以及其他元数据信息一起存储在文件的末尾。

文件存储格式示例:

1
4-byte magic number "PAR1"

整个文件(表)有 N 个列,划分成了 M 个行组,每个行组都有所有列的一个 Chunk 和其元数据信息。文件的元数据信息存储在数据之后,包含了所有列块元数据信息的起始位置。读取的时候首先从文件末尾读取文件元数据信息,再在其中找到感兴趣的 Column Chunk 信息,并依次读取。文件元数据信息放在文件最后是为了方便数据依序一次性写入。

具体的存储格式展示图:

图片

1.3.4 元数据信息

Parquet 总共有 3 种类型的元数据:文件元数据、列(块)元数据和 page header 元数据。所有元数据都采用 thrift 协议存储。具体信息如下所示:

图片

1.3.5 Parquet 数据类型

在实现层级上,Parquet 只保留了最精简的部分数据类型,以方便存储和读写。在其上有逻辑类型(Logical Types)以供扩展,比如:逻辑类型 strings 就映射为带有 UTF8 标识的二进制 byte arrays 进行存储。

Types:

1
BOOLEAN: 1 bit boolean

逻辑类型的更多说明请参考:

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md

1.3.6 Encoding

数据编码的实现大部分和原理部分所阐述的一致,这里不再重复说明,更多细节可参考:https://github.com/apache/parquet-format/blob/master/Encodings.md

1.3.7 Column chunks 存储

Column chunks 由一个个 Pages 组成,Reader 在读取的时候可以根据 page header 信息跳过不感兴趣的页面。page header 中还存储着页面数据编码和压缩的信息。

1.3.8 错误情况处理

如果文件元数据损坏,则整个文件将丢失。如果列元数据损坏,则该列块将丢失(但其他行组中该列的列块还可以使用)。如果 page header 损坏,则该列块中的剩余页面都将丢失。如果页面中的数据损坏,则该页面将丢失。较小的文件行组配置,可以更有效地抵抗损坏。

1.3.9 推荐配置

行组大小(Row group size): 更大的行组允许更大的列块,这使得可以执行更大的顺序 IO。不过更大的行组需要更大的写缓存。Parquet 建议使用较大的行组(512MB-1GB)。此外由于可能需要读取整个行组,因此最好一个行组能完全适配一个 HDFS Block。因此,HDFS 块大小也需要相应的设置更大。一个较优的读取配置为:行组大小 1GB,HDFS 块大小 1GB,每个 HDFS 文件对应 1 个 HDFS 块。

**数据页大小(Data page size):**数据页应视为不可分割的,因此较小的数据页可实现更细粒度的读取(例如单行查找)。但较大的页面可以减少空间的开销(减少 page header 数量)和潜在的较少的解析开销(处理 headers)。Parquet 建议的页面大小为 8KB。

1.4 Parquet分析工具

https://github.com/hangxie/parquet-tools/blob/main/USAGE.md#brew-install

参考资料:

[1]. Dremel: Interactive Analysis of WebScale Datasets
[2]. Dremel made simple with Parquet
[3]. 经典论文翻译导读之《Dremel: Interactive Analysis of WebScale Datasets》
[4]. 处理海量数据:列式存储综述(存储篇)
[5]. https://parquet.apache.org/documentation/latest/
[6]. https://blog.csdn.net/com360/article/details/13774489
[7]. 详解Parquet文件格式 - 深潜