工作流调度 Azkaban 简介和安装

工作流调度器azkaban

https://azkaban.readthedocs.io/en/latest/getStarted.html

一、课前准备

  1. 安装VMware虚拟化软件
  2. 安装CentOS 7虚拟机3个
  3. 安装3节点的hadoop集群
  4. 安装了hive
  5. 安装了zookeeper集群
  6. 安装了hbase集群

二、课堂主题

  1. azkaban架构
  2. azkaban运行模式
  3. azkaban安装部署
  4. azkaban使用

三、课堂目标

  1. 理解azkanban架构
  2. 完成azkaban安装部署
  3. 学会azkaban各种使用方式

四、知识要点

1. 概述

1. 为什么需要工作流调度系统

  • 一个完整的数据分析系统通常都是由大量任务单元组成:
    • shell脚本程序,java程序,mapreduce程序、hive脚本等
    • 各任务单元之间存在时间先后及前后依赖关系
    • 为了很好地组织起这样的复杂执行计划,需要一个工作流调度系统来调度执行;
  • 例如,我们可能有这样一个需求,某个业务系统每天产生20G原始数据,我们每天都要对其进行处理,处理步骤如下所示:
    • 通过Hadoop先将原始数据同步到HDFS上;
    • 借助MapReduce计算框架对原始数据进行转换,生成的数据以分区表的形式存储到多张Hive表中;
    • 需要对Hive中多个表的数据进行JOIN处理,得到一个明细数据Hive宽表;
    • 将明细数据进行各种统计分析,得到结果报表信息;
    • 需要将统计分析得到的结果数据同步到业务系统中,供业务调用使用。

2. 工作流调度实现方式

  • 简单的任务调度:直接使用linux的crontab来定义;

  • 复杂的任务调度:开发调度平台或使用现成的开源调度系统,比如ooize、azkaban、airflow、dophinschedule等

2. Azkaban介绍

  • Azkaban是由Linkedin开源的一个批量工作流任务调度器。用于在一个工作流(work flow)内以一个特定的顺序运行一组工作和流程。

  • Azkaban定义了一种KV格式文件(properties)来建立任务之间的依赖关系,并提供一个易于使用的web用户界面维护和跟踪你的工作流。

  • 它有如下功能特点:

    • 提供功能清晰、简单易用的web UI界面

    • 方便上传工作流

    • 调度工作流

    • 能够杀死并重新启动工作流

    • 工作流和任务的日志记录和审计

    • 提供job配置文件快速建立任务和任务之间的关系

    • 提供模块化的可插拔机制,原生支持command、java、hive、hadoop

    • 安全性高:认证/授权(权限的工作)

    • 提供分布式的多个执行服务器executor

    • 提供conditional workflow工作流

3. azkaban的基本架构

img

  • Azkaban由三部分构成
    • 1、Azkaban Web Server
      提供了Web UI,是azkaban的主要管理者,包括 project 的管理,认证,调度,对工作流执行过程的监控等。
    • 2、Azkaban Executor Server
      负责具体的工作流和任务的调度提交
    • 3、Mysql
      用于保存项目、日志或者执行计划之类的信息

4. Azkaban架构的三种运行模式

1. solo server mode(单机模式)

  • solo server mode是azkaban的一个独立的实例

  • 易于安装:不需要安装mysql,它内置了H2数据库,作为它的底层持久化存储

  • 易于开始使用:管理服务器web server和执行服务器execute server都在一个进程中运行,任务量不大项目可以采用此模式

  • 包含azkaban所有的功能

  • 有兴趣的同学,可以参考官网文档

2. two server mode

  • web server 和 executor server运行在不同的进程
  • 数据库为mysql,管理服务器和执行服务器在不同进程
  • 这种模式下,管理服务器和执行服务器互不影响。

3. multiple executor mode

  • web server 和 executor server运行在不同的进程,executor server有多个
  • 该模式下,执行服务器和管理服务器分别部署,且执行服务器可以在不同服务器上有多个运行的实例。

5. Azkaban安装部署

1. 编译azkaban

建议:直接使用老师编译出来的安装包进行安装

所以“编译azkaban”这个步骤,可以做个了解即可,不用亲自动手编译

1、下载源码包

这里选用azkaban4.0.0这个版本的源码进行重新编译,编译完成之后得到我们需要的安装包,然后进行安装

注意:

  • 使用Gradle编译azkaban源码
  • 需要使用jdk1.8或更高的版本来进行编译
  • 先获得azkaban源码
  • 浏览器访问地址https://github.com/azkaban/azkaban/releases

下载源码包

image-20210318164312734

  • 或者使用命令下载
cd /export/softwares/
wget https://github.com/azkaban/azkaban/archive/4.0.0.tar.gz
2、修改build.gradle
  • azkaban-4.0.0.tar.gz源码包上传到编译源码的虚拟机的/export/softwares目录,然后解压并编译

    提示:最好有一个虚拟机,专门用于编译各种框架的源码

# 切换到root用户
su root
cd /export/softwares/
tar -xzvf azkaban-4.0.0.tar.gz -C /export/servers/
cd /export/servers/azkaban-4.0.0/

使用gradle进行编译源码,此过程中

  • 需要去maven的仓库中下载各种jar包等文件,为了提高下载的速度,可以配置成从国内的maven仓库下载文件
  • 方法如下
vim /export/servers/azkaban-4.0.0/build.gradle

在如下2个位置添加maven仓库url

  maven{ url 'http://maven.aliyun.com/nexus/content/groups/public/'}
  maven{ url 'http://maven.oschina.net/content/groups/public/'}

位置①

image-20210318170454256

位置②

image-20210318170325962

3、开始编译
cd /export/servers/azkaban-4.0.0
yum -y install git
yum -y install gcc-c++
./gradlew build installDist -x test

#开始下载,控制台会打印如下类似的日志
Downloading https://services.gradle.org/distributions/gradle-4.6-all.zip
..........
  • 编译成功

image-20210323145150033

编译之后,获得安装包如下

4、azkaban-exec-server

编译完成之后得到我们需要的安装包在以下目录下即可获取得到azkaban-exec-server存放目录

cd /export/servers/azkaban-4.0.0/azkaban-exec-server/build/distributions/
ll

image-20210318171406141

5、azkaban-web-server

azkaban-web-server存放目录

cd /export/servers/azkaban-4.0.0/azkaban-web-server/build/distributions/
ll

image-20210318171623513

6、azkaban-solo-server

azkaban-solo-server存放目录

cd /export/servers/azkaban-4.0.0/azkaban-solo-server/build/distributions/
ll

image-20210318171757237

7、execute-as-user.c

azkaban two server模式下需要的C程序在这个路径下面

cd /export/servers/azkaban-4.0.0/az-exec-util/src/main/c
ll

image-20210318171924253

8、数据库脚本文件

数据库脚本文件在这个路径下面

cd /export/servers/azkaban-4.0.0/azkaban-db/build/sql/
ll

image-20210318172146694

如果不想手动编译可以直接使用我这里编译好的文件,百度网盘链接:https://pan.baidu.com/s/14owbePcSBXBW_bwyLTVVNQ
提取码:NIIT

2. multiple executor模式安装

前提:某节点已经安装mysql,此文档以node03已经安装mysql为例

若没有特殊说明,所有操作都是使用hadoop普通用户操作

在node03节点操作

1. 确认所需软件:

Azkaban Web服务安装包

azkaban-web-server-0.1.0-SNAPSHOT.tar.gz

Azkaban执行服务安装包

azkaban-exec-server-0.1.0-SNAPSHOT.tar.gz

编译之后的sql脚本

create-all-sql-0.1.0-SNAPSHOT.sql

C程序文件脚本

execute-as-user.c

将以上4个文件上传到node03的/export/softwares目录

2. 数据库准备

进入mysql的客户端执行以下命令

[hadoop@node03 ~]$ mysql -uroot -p123456

mysql中执行以下命令:

-- 设置密码的验证强度等级
set global validate_password_policy=LOW; 
set global validate_password_length=6;

-- 创建数据库azkaban,用于存储使用azkaban框架过程中产生的数据
CREATE DATABASE azkaban;
CREATE USER 'azkaban'@'%' IDENTIFIED BY 'azkaban';  
GRANT SELECT,INSERT,UPDATE,DELETE ON azkaban.* to 'azkaban'@'%' identified by 'azkaban' WITH GRANT OPTION;

flush privileges;
use azkaban; 
source /export/softwares/create-all-sql-0.1.0-SNAPSHOT.sql;
exit;

image-20210319102927859

更改 MySQL 包大小;防止 Azkaban 连接 MySQL 阻塞

hadoop@node03 ~]$ sudo vim /etc/my.cnf

在文件末尾增加如下内容;然后保存、退出

max_allowed_packet=1024M

image-20210319163143002

重启mysql服务

[hadoop@node03 sbin]$ sudo /sbin/service mysqld restart
# 输出如下日志
Redirecting to /bin/systemctl restart mysqld.service
3. 解压软件安装包

解压azkaban-web-server

[hadoop@node03 ~]$ cd
[hadoop@node03 ~]$ cd /export/softwares/
[hadoop@node03 softwares]$ tar -zxvf azkaban-web-server-0.1.0-SNAPSHOT.tar.gz -C /export/servers
[hadoop@node03 softwares]$ cd /export/servers
[hadoop@node03 servers]$ mv azkaban-web-server-0.1.0-SNAPSHOT/ azkaban-web-server-4.0.0

解压azkaban-exec-server

[hadoop@node03 servers]$ cd /export/softwares/
[hadoop@node03 softwares]$ tar -zxvf azkaban-exec-server-0.1.0-SNAPSHOT.tar.gz -C /export/servers
[hadoop@node03 softwares]$ cd /export/servers
[hadoop@node03 servers]$ mv azkaban-exec-server-0.1.0-SNAPSHOT/ azkaban-exec-server-4.0.0
4. 安装SSL安全认证

安装ssl安全认证,允许我们使用https的方式访问我们的azkaban的web服务;

密码一定要一个个的字母输入,或者粘贴也行

[hadoop@node03 servers]$ cd /export/servers/azkaban-web-server-4.0.0
[hadoop@node03 azkaban-web-server-4.0.0]$ keytool -keystore keystore -alias jetty -genkeypair -keyalg RSA

具体输入如下:

  • 密码都是azkaban
  • 显示[Unknown]:直接敲回车
  • 显示[no]:输入yes,回车
  • jetty密码:azkaban

如果你的linux是中文环境

image-20210319161248569

如果是英文环境,上图中的“是”用“yes”代替
发现目录中多出一个秘钥文件keystore

[hadoop@node03 azkaban-web-server-4.0.0]$ pwd
/export/servers/azkaban-web-server-4.0.0
[hadoop@node03 azkaban-web-server-4.0.0]$ ls
bin  conf  keystore  lib  web

补充:keytool的用法

# keytool中有很多的cmd及option;查看keytool用法
[hadoop@node03 azkaban-web-server-4.0.0]$ man keytool

# 提示:搜索关键字,如-genkeypair

OPTION DEFAULTS
       The following examples show the defaults for various option values.
       -alias "mykey"
       -keystore <the file named .keystore in the user's home directory>
       -genkeypair 生成一个秘钥对(公钥、私钥)
       -keyalg RSA 指定使用某算法生成秘钥对
5. azkaban web server安装
1、修改azkaban-web-server

修改azkaban-web-server的配置文件

[hadoop@node03 azkaban-web-server-4.0.0]$ cd /export/servers/azkaban-web-server-4.0.0/conf
[hadoop@node03 conf]$ vim azkaban.properties

修改文件中的如下属性

# Azkaban Personalization Settings
azkaban.name=Azkaban
azkaban.label=My Azkaban
...

default.timezone.id=Asia/Shanghai
...

# Azkaban Jetty server properties.
jetty.use.ssl=true
...

# 新增内容
jetty.ssl.port=8443
jetty.keystore=/export/servers/azkaban-web-server-4.0.0/keystore
jetty.password=azkaban
jetty.keypassword=azkaban
jetty.truststore=/export/servers/azkaban-web-server-4.0.0/keystore
jetty.trustpassword=azkaban
...

mysql.host=node03
...

azkaban.executorselector.filters=StaticRemainingFlowSize,CpuStatus

说明:

  • StaticRemainingFlowSize:正在排队的任务数;
  • CpuStatus: CPU 占用情况
  • MinimumFreeMemory:内存占用情况。 测试环境, 必须将 MinimumFreeMemory 删除掉,否则它会认为集群资源不够,不执行。

需要修改的项目如下图

image-20210319172111245

2、设置azkaban用户

文件/export/servers/azkaban-web-server-4.0.0/conf/azkaban.properties中默认有如下属性

user.manager.class=azkaban.user.XmlUserManager
user.manager.xml.file=conf/azkaban-users.xml

上边的XmlUserManager是azkaban内置的UserManager用户管理器

  • 当启动web server时,XmlUserManager读取上边配置文件azkaban.properties,然后解析azkaban-users.xml
  • 此xml文件中用来配置azkaban的用户、组、角色
  • 接下来编辑此xml文件,用来配置azkaban用户
  • 添加用户hadoop
[hadoop@node03 conf]$ pwd
/export/servers/azkaban-web-server-4.0.0/conf
[hadoop@node03 conf]$ vim azkaban-users.xml

内容如下(注意:在向内容拷贝到xml文件时,注意格式缩进)

  <user password="kkb123" roles="myread" username="kkbread"/>
  <user password="kkb123" roles="mywrite" username="kkbwrite"/>
  <user password="kkb123" roles="admin" username="kkbadmin"/>
  <user password="kkb123" roles="myread, mywrite" username="kkbrwe" groups="groupx"/>

  <group name="groupx" roles="myexe"/>

  <role name="myread" permissions="READ"/>
  <role name="mywrite" permissions="WRITE"/>
  <role name="myexe" permissions="EXECUTE"/>  

image-20210319190544587

此文件格式解析,可以参考官网文档

user格式

Attributes Values Required?
username 登录用户名 The login username. 必须有
password 密码 The login password. 必须有
roles role角色,如果是多个角色的话,中间逗号分隔
Comma delimited list of roles that this user has.
非必须
groups 用户所属组,如果是多个组,中间逗号分隔
Comma delimited list of groups that the users belongs to.
非必须
proxy 代理 Comma delimited list of proxy users that this users can give to a project 非必须

group

Attributes Values Required?
name 组名 The group name 必须有
roles role角色,如果是多个角色的话,中间逗号分隔
Comma delimited list of roles that this user has.
非必须

role

Attributes Values Required?
name role角色名称 The role name 必须有
permissions 权限,如果是多个权限的话,中间逗号分隔
Comma delimited list global permissions for the role
必须有

可选的权限有

Permissions Values
ADMIN 管理员权限,拥有azkaban中所有的权限
Grants all access to everything in Azkaban.
READ 对每个project有只读权限
Gives users read only access to every project and their logs
WRITE 允许用户上传文件、修改job的properties、删除project
Allows users to upload files, change job properties or remove any project
EXECUTE 允许用户执行任何flow
Allows users to trigger the execution of any flow
SCHEDULE 允许用户给任意flow添加或移除指定的调度
Users can add or remove schedules for any flows
CREATEPROJECTS 如果创建project功能被锁死,有此权限的用户拥有创建project的权限
Allows users to create new projects if project creation is locked down

官网配置

6. azkaban executor server 安装
第一步:修改azkaban-exex-server配置文件

修改azkaban-exec-server的配置文件

[hadoop@node03 ~]$ cd /export/servers/azkaban-exec-server-4.0.0/conf
[hadoop@node03 conf]$ vim azkaban.properties
# Azkaban Personalization Settings
azkaban.name=Azkaban
azkaban.label=My Azkaban
...

default.timezone.id=Asia/Shanghai
...

jetty.use.ssl=true
......

# 新增内容 添加如下5行内容
jetty.keystore=/export/servers/azkaban-web-server-4.0.0/keystore
jetty.password=azkaban
jetty.keypassword=azkaban
jetty.truststore=/export/servers/azkaban-web-server-4.0.0/keystore
jetty.trustpassword=azkaban
...

# Where the Azkaban web server is located
azkaban.webserver.url=https://node03:8443
...

mysql.host=node03
...
第二步:添加插件

将我们编译后的C文件execute-as-user.c上传或拷贝到/export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes

[hadoop@node03 conf]$ cp /export/softwares/execute-as-user.c /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes/

然后执行以下命令生成execute-as-user

# 在线安装gcc-c++
[hadoop@node03 conf]$ sudo yum -y install gcc-c++
[hadoop@node03 conf]$ cd /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes
[hadoop@node03 jobtypes]$ gcc execute-as-user.c -o execute-as-user
# 添加root特权
[hadoop@node03 jobtypes]$ sudo chown root execute-as-user
[sudo] hadoop 的密码:
[hadoop@node03 jobtypes]$ sudo chmod 6050 execute-as-user

image-20210321122703300

备注:数字表示的权限为6050,那对应的字母表示是什么呢?

  • 第一数字6表示的是特殊权限6 = 4 + 2,即同时设置了SUID和SGID
  • 第二个数字0表示的所有者权限为0,字母表示为---,因为没有设置x权限, 特殊权限表示为大写字母S
  • 第三个数字为5 =4 + 1,设置了用户组权限为读(r)和执行(x),特殊权限表示为小写字母s
  • 第四个数字为0,字母表示为---

其实在UNIX的实现中,文件权限用12个二进制位表示,如果该位置上的值是1,表示有相应的权限:|

11 10 9  8  7  6  5  4  3  2  1  0
S  G  T  r  w  x  r  w  x  r  w  x

第11位为SUID(Set User ID)位,第10位为SGID(Set Group ID)位,第9位为sticky位,第8-0位对应于上面的三组rwx位。

-rwsr-xr-x的值为: 表示所属用户有特殊权限SUID

1 0 0 1 1 1 1 0 1 1 0 1

-rw-r-Sr--的值为: 表示所属用户组有特殊权限SGID

0 1 0 1 1 0 1 0 0 1 0 0

给文件加SUIDSUID的命令如下:

chmod u+s filename 设置SUID

chmod u-s filename 去掉SUID设置

chmod g+s filename 设置SGID

chmod g-s filename 去掉SGID设置

另外一种方法是chmod 接4位八进制表示法。

比如 chmod 60506050对应的2进制110 000 101 000

对应的是权限是 SG --- rx ---

就是用户组有读写权限,并且同时有用户和用户组的特殊权限

但是使用命令ll列出文件的权限表示是这样

--S r-s --- 这是因为其中S表示设置SUID,s表示设置SGID同时用户组有执行权限

SUIDSGID的作用:

随意设置会威胁到系统安全因此终端上显示成红色背景,特殊权限只在执行普通文件时生效。

比如本来是只有root用户才能执行的命令,加了SUID后,普通用户就可以像root一样用这个命令,权限提升了。上面是对于文件来说的,对于目录也差不多!

目录的S属性使得在该目录下创建的任何文件及子目录属于该目录所拥有的组,目录的T属性使得该目录的所有者及root才能删除该目录。还有对于sS,设置SUID/SGID需要有运行权限,否则用ls -l后就会看到S,证明你所设置的SUID/SGID没有起作用。

第三步:修改配置文件

修改配置文件

[hadoop@node03 jobtypes]$ cd /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes
[hadoop@node03 jobtypes]$ vim commonprivate.properties

增加或修改如下内容

execute.as.user=true
azkaban.native.lib=/export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes
azkaban.group.name=myazkaban
memCheck.enabled=false
azkaban.native.lib=false

遇到报错Missing required property “azkaban.native.lib”的解决办法:/plugins/jobtypes 目录下修改commonprivate.properties配置文件,内容中添加:azkaban.native.lib=false。然后重启启动exec服务,并激活executor。

将exec拷贝到另外两个节点,并修改所属用户

[hadoop@node03 servers]$ cd /export/servers/
[hadoop@node03 servers]$ scp -r azkaban-exec-server-4.0.0/ root@node01:$PWD
[hadoop@node03 servers]$ scp -r azkaban-exec-server-4.0.0/ root@node02:$PWD

node01节点

[root@node01 ~]# su root
[root@node01 servers]# cd /export/servers/
[root@node01 servers]# chown -R hadoop:hadoop azkaban-exec-server-4.0.0/
[root@node01 servers]# su hadoop
[hadoop@node01 jobtypes]$ cd /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes
# 添加root特权
[hadoop@node01 jobtypes]$ sudo chown root execute-as-user
[sudo] hadoop 的密码:
[hadoop@node01 jobtypes]$ sudo chmod 6050 execute-as-user

node02节点

[hadoop@node02 ~]$ su root
密码:
[root@node02 servers]# cd /export/servers/
[root@node02 servers]# chown -R hadoop:hadoop azkaban-exec-server-4.0.0/

[root@node02 servers]# su hadoop
[hadoop@node02 jobtypes]$ cd /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes
# 添加root特权
[hadoop@node02 jobtypes]$ sudo chown root execute-as-user
[sudo] hadoop 的密码:
[hadoop@node02 jobtypes]$ sudo chmod 6050 execute-as-user
第四步:添加用户、用户组

web server

  • node03节点/export/servers/azkaban-web-server-4.0.0/conf/azkaban-users.xml中,添加了许多用户,如kkbreadkkbrwe等等
  • 登录web server时,就用这里指定的用户名、密码登录

image-20210323154047477

exec server

  • /export/servers/azkaban-exec-server-4.0.0/plugins/jobtypes/commonprivate.properties
  • 属性execute.as.user=true表示,azkaban的登录用户,同时作为linux服务器的系统用户
  • 属性azkaban.group.name=myazkaban表示,这些azkaban用户(linux系统用户)都属于myazkaban用户组
  • 所以需要在3个exec服务器中,创建用户kkbrwe、用户组myazkaban
  • 并且用户kkbrwe属于用户组myazkaban
  • 为了解决权限问题,同时讲用户kkbrwe添加附属组hadoop
  • 此处以kkbrwe用户为例(如果使用其他用户,按照此方式创建即可)具体命令如下
  • 三节点hadoop用户已经有sudoers权限
  • 三节点 node01、node02、node03都执行如下命令
sudo groupadd myazkaban
[sudo] hadoop 的密码:
sudo useradd -g myazkaban kkbrwe
sudo passwd kkbrwe
新的 密码:123456
重新输入新的密码:
# 将kkbrwe添加附加用户组
sudo usermod -a -G hadoop kkbrwe
# 查看用户kkbrwe
sudo id kkbrwe
uid=1001(kkbrwe) gid=1001(myazkaban) 组=1001(myazkaban),1000(hadoop)

exec服务器在执行flow时,会在/export/servers/azkaban-exec-server-4.0.0目录创建目录executions,为了解决权限问题,3台节点都需要做如下操作

node01节点

[hadoop@node01 ~]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node01 azkaban-exec-server-4.0.0]$ mkdir executions
[hadoop@node01 azkaban-exec-server-4.0.0]$ sudo chown :myazkaban executions/

node02节点

[hadoop@node02 ~]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node02 azkaban-exec-server-4.0.0]$ mkdir executions
[hadoop@node02 azkaban-exec-server-4.0.0]$ sudo chown :myazkaban executions/

node03节点

[hadoop@node03 ~]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node03 azkaban-exec-server-4.0.0]$ mkdir executions
[hadoop@node03 azkaban-exec-server-4.0.0]$ sudo chown :myazkaban executions/

说明::myazkaban对应exec的commonprivate.propertiesazkaban.group.name=myazkaban属性

7. 启动服务
第一步:启动azkaban exec server
  • node01节点上
[root@node01 servers]# cd /export/servers/
[root@node01 servers]# su hadoop
[hadoop@node01 servers]$ cd azkaban-exec-server-4.0.0/
[hadoop@node01 azkaban-exec-server-4.0.0]$ bin/start-exec.sh
[hadoop@node01 azkaban-exec-server-4.0.0]$ jps

image-20210321163632014

关闭exec server: bin/shutdown-exec.sh

第二步:激活我们的exec-server
  • 每次启动exec都需要激活

  • node01机器下执行以下命令

[hadoop@node01 azkaban-exec-server-4.0.0]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node01 azkaban-exec-server-4.0.0]$ curl -G "node01:$(<./executor.port)/executor?action=activate" && echo
{"status":"success"}

image-20210321164134373

在node02上启动exec server并激活(同上)

[root@node02 servers]# cd /export/servers/
[root@node02 servers]# su hadoop
[hadoop@node02 servers]$ cd azkaban-exec-server-4.0.0/
[hadoop@node02 azkaban-exec-server-4.0.0]$ bin/start-exec.sh
[hadoop@node02 azkaban-exec-server-4.0.0]$ jps
[hadoop@node02 azkaban-exec-server-4.0.0]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node02 azkaban-exec-server-4.0.0]$ curl -G "node02:$(<./executor.port)/executor?action=activate" && echo
{"status":"success"}

在node03上启动exec server并激活(同上)

[hadoop@node03 servers]$ cd /export/servers/
# 如果当前用户不是hadoop,则切换到hadoop;否则,可以省略此命令
[hadoop@node03 servers]$ su hadoop
密码:
[hadoop@node03 servers]$ cd azkaban-exec-server-4.0.0/
[hadoop@node03 azkaban-exec-server-4.0.0]$ bin/start-exec.sh
[hadoop@node03 azkaban-exec-server-4.0.0]$ jps
17203 AzkabanExecutorServer
17224 Jps
[hadoop@node03 azkaban-exec-server-4.0.0]$ cd /export/servers/azkaban-exec-server-4.0.0
[hadoop@node03 azkaban-exec-server-4.0.0]$ curl -G "node03:$(<./executor.port)/executor?action=activate" && echo
{"status":"success"}

三个节点的exec server都启动后,可以去mysql中确认下,node03执行命令

[hadoop@node03 azkaban-web-server-4.0.0]$ mysql -uroot -p
Enter password:
mysql> use azkaban;
mysql> show tables;
mysql> select * from executors;
+----+--------+-------+--------+
| id | host   | port  | active |
+----+--------+-------+--------+
|  1 | node01 | 39157 |      1 |
|  2 | node02 | 41613 |      1 |
|  3 | node03 | 34589 |      1 |
+----+--------+-------+--------+
3 rows in set (0.00 sec)

发现,确实有3个exec server,active=1,表示已激活

也可以在MySQL中使用update语句设置active为1完成激活

第三步:启动azkaban-web-server
  • node03节点
[hadoop@node03 azkaban-exec-server-4.0.0]$ cd /export/servers/azkaban-web-server-4.0.0/
[hadoop@node03 azkaban-web-server-4.0.0]$ bin/start-web.sh
[hadoop@node03 azkaban-web-server-4.0.0]$ jps
17298 AzkabanWebServer
17203 AzkabanExecutorServer
17322 Jps

关闭web server命令:bin/shutdown-web.sh

注意,无论启动一定要在server根目录下使用相对命令的方式执行,如 bin/cmd,这是因为启动脚本中使用了相对路径读取权限相关的xml文件。

image-20210321164822923

宿主机浏览器访问地址:https://node03:8443

前提:宿主机的hosts文件中配置了node03 ip地址与主机名node03的映射,类似

image-20210321170137537

注意一定要使用https://协议访,问若浏览器界面出现类似情况,按图操作即可

image-20200803131019862

image-20200803131037682

image-20210322141415494

用户名和密码可以是node03的文件/export/servers/azkaban-web-server-4.0.0/conf中指定的用户名及密码,如下

比如此处使用用户kkbrwe,密码kkb123

image-20210321170311511

登录后,进入界面

image-20210322141507162

8. 修改linux的时区问题

之前在安装虚拟机时,已经设置时区为“亚洲/上海”所以不用担心时区问题,不需要修改时区

注:先配置好服务器节点上的时区

但是如果你的时区不是“亚洲/上海”,那么需要修改成此时区

确认时区;CST +0800表示时区是东八区(“亚洲/上海”)

[root@localhost azkaban-master]# date +"%Z %z"
CST +0800
  • 1、先生成时区配置文件Asia/Shanghai,用交互式命令 tzselect 即可
  • 2、拷贝该时区文件,覆盖系统本地时区配置
cp /usr/share/zoneinfo/Asia/Shanghai /etc/localtime

Views: 49

大数据岗位需求情况分析(二)结果导出和可视化

Sqoop导出Hive表到MySQL中

前提

  • DFS和Yarn需保持运行状态
  • MySQL服务处于运行状态

创建对应的MySQL表

根据Hive中用于存放分析结果的四个表,也同样在MySQL中创建具有相同表结构的四个表

使用Sqoop命令导出

语法

sqoop export \
--connect "jdbc:mysql://hadoop100:3306/job?useSSL=false&characterEncoding=utf-8"  \
--username root --password niit1234 \
--table <mysql_table_name> \
--export-dir /user/hive/warehouse/job.db/<hive_table_name>/data_date=<partition_value> \
--input-fields-terminated-by "\001";
第一个表

建表

mysql> create table job_count(process_date date,total_job_count int);

导出

sqoop export \
--connect "jdbc:mysql://hadoop100:3306/job?useSSL=false&characterEncoding=utf-8"  \
--username root \
--password niit1234 \
--table job_count \
--export-dir /user/hive/warehouse/job.db/job_count/data_date=2021-02-28 \
--input-fields-terminated-by "\001"
第二个表

建表

mysql> create table job_city_count(process_date date,city varchar(20), total_job_count int);

导出

sqoop export \
--connect "jdbc:mysql://hadoop100:3306/job?useSSL=false&characterEncoding=utf-8"  \
--username root \
--password niit1234 \
--table job_city_count \
--export-dir /user/hive/warehouse/job.db/job_city_count/data_date=2021-02-28 \
--input-fields-terminated-by "\001"
第三个表

建表

 create table job_city_salary(process_date date,city varchar(20), job_name varchar(50), salary_per_month int);

导出

sqoop export \
--connect "jdbc:mysql://hadoop100:3306/job?useSSL=false&characterEncoding=utf-8"  \
--username root \
--password niit1234 \
--table job_city_salary \
--export-dir "/user/hive/warehouse/job.db/job_city_salary/data_date=2021-02-28" \
--input-fields-terminated-by "\001"
第四个表

建表

create table job_tag(process_date date,job_tag varchar(50), tag_count int);

导出

sqoop export \
--connect "jdbc:mysql://hadoop100:3306/job?useSSL=false&characterEncoding=utf-8"  \
--username root \
--password niit1234 \
--table job_tag \
--export-dir "/user/hive/warehouse/job.db/job_tag/data_date=2021-02-28" \
--input-fields-terminated-by "\001"

数据可视化展示

superset是由Airbnb(知名在线短租赁公司)开源的数据分析与可视化平台(曾用名Caravel、Panoramix),该工具主要特点是可自助分析、自定义仪表盘、分析结果可视化(导出)、用户/角色权限控制,还集成了一个SQL编辑器,可以进行SQL编辑查询对结果集进行保存可视化等。

SUPERSET的基本介绍与安装参考这篇文章

启动superset后,打开浏览器访问http://hadoop100:8787

输入用户名(admin)和密码(admin)登陆即可使用superset进行数据可视化展示。

创建数据库连接

file
具体配置为
file

创建数据集

file

创建图表

根据数据集创建需要展示的图标(Chart)

表1

爬取的总岗位数
file

表2

不同城市提供的大数据相关岗位数量比较

file

表3

不同城市提供的大数据相关岗位的薪资倒序排列 - 取TopN
file

表4

岗位标签做成词云统计

file

创建仪表盘

仪表盘可以将需要展示的所有图标布局到一起。

file
布局后

file

将仪表盘设置为实时更新

如果是需要实时更新数据的表,可以设置同步间隔时间

file

Views: 99

Flume Sinks

Flume Sinks 类型有很多,这里只挑出一些我们常用的Sink.

  1. HDFS Sink
  2. Hive Sink
  3. Logger Sink
  4. Avro Sink
  5. HBase Sinks
  6. Kafka Sink
  7. HTTP Sink
  8. File Roll Sink
  9. NULL sink
  10. Custom SInk

HDFS Sink

这个Sink将Event写入Hadoop分布式文件系统(也就是HDFS)。 目前支持创建文本和序列文件。 它支持两种文件类型的压缩。 可以根据写入的时间、文件大小或Event数量定期滚动文件(关闭当前文件并创建新文件)。 它还可以根据Event自带的时间戳或系统时间等属性对数据进行分区。 存储文件的HDFS目录路径可以使用格式转义符,会由HDFS Sink进行动态地替换,以生成用于存储Event的目录或文件名。

使用此Sink需要安装hadoop, 以便Flume可以使用Hadoop的客户端与HDFS集群进行通信。

注意,%[localhost], %[IP] 和 %[FQDN]这三个转义符实际上都是用java的API来获取的,在一些网络环境下可能会获取失败。

正在打开的文件会在名称末尾加上“.tmp”的后缀。文件关闭后,会自动删除此扩展名。这样容易排除目录中的那些已完成的文件。 必需的参数已用 粗体 标明。

属性名默认值解释
channel与 Sink 连接的 channel
type组件类型,这个是: hdfs
hdfs.pathHDFS目录路径(例如:hdfs://namenode/flume/webdata/)
hdfs.filePrefixFlumeDataFlume在HDFS文件夹下创建新文件的固定前缀
hdfs.fileSuffixFlume在HDFS文件夹下创建新文件的后缀(比如:.avro,注意这个“.”不会自动添加,需要显式配置)
hdfs.inUsePrefixFlume正在写入的临时文件前缀,默认没有
hdfs.inUseSuffix.tmpFlume正在写入的临时文件后缀
hdfs.emptyInUseSuffixfalse如果设置为 false 上面的 hdfs.inUseSuffix 参数在写入文件时会生效,并且写入完成后会在目标文件上移除 hdfs.inUseSuffix 配置的后缀。如果设置为 true 则上面的 hdfs.inUseSuffix 参数会被忽略,写文件时不会带任何后缀
hdfs.rollInterval30当前文件写入达到该值时间后触发滚动创建新文件(0表示不按照时间来分割文件),单位:秒
hdfs.rollSize1024当前文件写入达到该大小后触发滚动创建新文件(0表示不根据文件大小来分割文件),单位:字节
hdfs.rollCount10当前文件写入Event达到该数量后触发滚动创建新文件(0表示不根据 Event 数量来分割文件)
hdfs.idleTimeout0关闭非活动文件的超时时间(0表示禁用自动关闭文件),单位:秒
hdfs.batchSize100向 HDFS 写入内容时每次批量操作的 Event 数量
hdfs.codeC压缩算法。可选值:gzip 、 bzip2 、 lzo 、 lzop 、 `snappy
hdfs.fileTypeSequenceFile文件格式,目前支持: SequenceFile 、 DataStream 、 CompressedStream 。 1. DataStream 不会压缩文件,不需要设置hdfs.codeC 2. CompressedStream 必须设置hdfs.codeC参数
hdfs.maxOpenFiles5000允许打开的最大文件数,如果超过这个数量,最先打开的文件会被关闭
hdfs.minBlockReplicas指定每个HDFS块的最小副本数。 如果未指定,则使用 classpath 中 Hadoop 的默认配置。
hdfs.writeFormatWritable文件写入格式。可选值: Text 、 Writable 。在使用 Flume 创建数据文件之前设置为 Text,否则 Apache Impala(孵化)或 Apache Hive 无法读取这些文件。
hdfs.threadsPoolSize10每个HDFS Sink实例操作HDFS IO时开启的线程数(open、write 等)
hdfs.rollTimerPoolSize1每个HDFS Sink实例调度定时文件滚动的线程数
hdfs.kerberosPrincipal用于安全访问 HDFS 的 Kerberos 用户主体
hdfs.kerberosKeytab用于安全访问 HDFS 的 Kerberos keytab 文件
hdfs.proxyUser 代理名
hdfs.roundfalse是否应将时间戳向下舍入(如果为true,则影响除 %t 之外的所有基于时间的转义符)
hdfs.roundValue1向下舍入(小于当前时间)的这个值的最高倍(单位取决于下面的 hdfs.roundUnit ) 例子:假设当前时间戳是18:32:01,hdfs.roundUnit = minute 如果roundValue=5,则时间戳会取为:18:30 如果roundValue=7,则时间戳会取为:18:28 如果roundValue=10,则时间戳会取为:18:30
hdfs.roundUnitsecond向下舍入的单位,可选值: second 、 minute 、 hour
hdfs.timeZoneLocal Time解析存储目录路径时候所使用的时区名,例如:America/Los_Angeles、Asia/Shanghai
hdfs.useLocalTimeStampfalse使用日期时间转义符时是否使用本地时间戳(而不是使用 Event header 中自带的时间戳)
hdfs.closeTries0开始尝试关闭文件时最大的重命名文件的尝试次数(因为打开的文件通常都有个.tmp的后缀,写入结束关闭文件时要重命名把后缀去掉)。如果设置为1,Sink在重命名失败(可能是因为 NameNode 或 DataNode 发生错误)后不会重试,这样就导致了这个文件会一直保持为打开状态,并且带着.tmp的后缀;如果设置为0,Sink会一直尝试重命名文件直到成功为止;关闭文件操作失败时这个文件可能仍然是打开状态,这种情况数据还是完整的不会丢失,只有在Flume重启后文件才会关闭。
hdfs.retryInterval180连续尝试关闭文件的时间间隔(秒)。 每次关闭操作都会调用多次 RPC 往返于 Namenode ,因此将此设置得太低会导致 Namenode 上产生大量负载。 如果设置为0或更小,则如果第一次尝试失败,将不会再尝试关闭文件,并且可能导致文件保持打开状态或扩展名为“.tmp”。
serializerTEXTEvent 转为文件使用的序列化器。其他可选值有: avro_event 或其他 EventSerializer.Builderinterface 接口的实现类的全限定类名。
serializer.* 根据上面 serializer 配置的类型来根据需要添加序列化器的参数

废弃的一些参数:

属性名默认值解释
hdfs.callTimeout10000允许HDFS操作文件的时间,比如:open、write、flush、close。如果HDFS操作超时次数增加,应该适当调高这个这个值。(毫秒)
配置范例:
a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = hdfs
a1.sinks.k1.channel = c1
a1.sinks.k1.hdfs.path = /flume/events/%y-%m-%d/%H%M/%S
a1.sinks.k1.hdfs.filePrefix = events-
a1.sinks.k1.hdfs.round = true
a1.sinks.k1.hdfs.roundValue = 10
a1.sinks.k1.hdfs.roundUnit = minute

上面的例子中时间戳会向前一个整10分钟取整。比如,一个 Event 的 header 中带的时间戳是11:54:34 AM, June 12, 2012,它会保存的 HDFS 路径就是/flume/events/2012-06-12/1150/00。

Hive Sink

此Sink将包含分隔文本或JSON数据的 Event 直接流式传输到 Hive表或分区上。 Event 使用 Hive事务进行写入, 一旦将一组 Event 提交给Hive,它们就会立即显示给Hive查询。 即将写入的目标分区既可以预先自己创建,也可以选择让 Flume 创建它们,如果没有的话。 写入的 Event 数据中的字段将映射到 Hive表中的相应列。

属性默认值解释
channel与 Sink 连接的 channel
type组件类型,这个是: hive
hive.metastoreHive metastore URI (eg thrift://a.b.com:9083 )
hive.databaseHive 数据库名
hive.tableHive表名
hive.partition逗号分隔的要写入的分区信息。 比如hive表的分区是(continent: string, country: string, time : string), 那么“Asia,India,2014-02-26-01-21”就表示数据会写入到continent=Asia,country=India,time=2014-02-26-01-21这个分区。
hive.txnsPerBatchAsk100Hive从Flume等客户端接收数据流会使用多次事务来操作,而不是只开启一个事务。这个参数指定处理每次请求所开启的事务数量。来自同一个批次中所有事务中的数据最终都在一个文件中。 Flume会向每个事务中写入 batchSize 个 Event,这个参数和 batchSize 一起控制着每个文件的大小,请注意,Hive最终会将这些文件压缩成一个更大的文件。
heartBeatInterval240发送到 Hive 的连续心跳检测间隔(秒),以防止未使用的事务过期。设置为0表示禁用心跳。
autoCreatePartitionstrueFlume 会自动创建必要的 Hive分区以进行流式传输
batchSize15000写入一个 Hive事务中最大的 Event 数量
maxOpenConnections500允许打开的最大连接数。如果超过此数量,则关闭最近最少使用的连接。
callTimeout10000Hive、HDFS I/O操作的超时时间(毫秒),比如:开启事务、写数据、提交事务、取消事务。
serializer 序列化器负责解析 Event 中的字段并把它们映射到 Hive表中的列,选择哪种序列化器取决于 Event 中的数据格式,支持的序列化器有:DELIMITED 和 JSON
roundfalse是否启用时间戳舍入机制
roundUnitminute舍入值的单位,可选值:second 、 minute 、 hour
roundValue1舍入到小于当前时间的最高倍数(使用 roundUnit 配置的单位) 例子1:roundUnit=second,roundValue=10,则14:31:18这个时间戳会被舍入到14:31:10; 例子2:roundUnit=second,roundValue=30,则14:31:18这个时间戳会被舍入到14:31:00,14:31:42这个时间戳会被舍入到14:31:30;
timeZoneLocal Time应用于解析分区中转义序列的时区名称,比如:America/Los_Angeles、Asia/Shanghai、Asia/Tokyo等
useLocalTimeStampfalse替换转义序列时是否使用本地时间戳(否则使用Event header中的timestamp )
下面介绍Hive Sink的两个序列化器:JSON :处理UTF8编码的 Json 格式(严格语法)Event,不需要配置。 JSON中的对象名称直接映射到Hive表中具有相同名称的列。 内部使用 org.apache.hive.hcatalog.data.JsonSerDe ,但独立于 Hive表的 Serde 。 此序列化程序需要安装 HCatalog。DELIMITED: 处理简单的分隔文本 Event。 内部使用 LazySimpleSerde,但独立于 Hive表的 Serde。

属性默认值解释
serializer.delimiter,(类型:字符串)传入数据中的字段分隔符。 要使用特殊字符,请用双引号括起来,例如“\t”
serializer.fieldnames从输入字段到Hive表中的列的映射。 指定为Hive表列名称的逗号分隔列表(无空格),按顺序标识输入字段。 要跳过字段,请保留未指定的列名称。 例如, ‘time,,ip,message’表示输入映射到hive表中的 time,ip 和 message 列的第1,第3和第4个字段。
serializer.serdeSeparatorCtrl-A(类型:字符)自定义底层序列化器的分隔符。如果 serializer.fieldnames 中的字段与 Hive表列的顺序相同,则 serializer.delimiter 与 serializer.serdeSeparator 相同, 并且 serializer.fieldnames 中的字段数小于或等于表的字段数量,可以提高效率,因为传入 Event 正文中的字段不需要重新排序以匹配 Hive表列的顺序。 对于’\t’这样的特殊字符使用单引号,要确保输入字段不包含此字符。 注意:如果 serializer.delimiter 是单个字符,最好将本参数也设置为相同的字符。

为了避免出现找不到类的异常,首先需添加依赖的jar包:

  1. 将hive/lib下面hive-hcatalog-core-3.1.2.jar拷贝或者软链接到flume/lib下
    1. $ cp /opt/pkg/hive/lib/hive-hcatalog-core-3.1.2.jar /opt/pkg/flume/lib/
  2. 接下来我们还需要这个依赖:hive-hcatalog-streaming-3.1.2.jar。这个hive目录是找不到的,需要单独从mvnrepository网站搜索下载,下载地址:https://mvnrepository.com/artifact/org.apache.hive.hcatalog/hive-hcatalog-streaming/3.1.2,下载后移动到flume/lib下即可。

使用Hive Sink 还有以下前提条件:

  • 需要开启事务支持
  • Hive表必须分区分桶
  • Hive表必须是Acid表(开启事务支持)

进入beeline临时开启事务支持(也可以修改配置文件永久开启,但不建议)


SET hive.support.concurrency = true;
SET hive.enforce.bucketing = true;
SET hive.exec.dynamic.partition.mode = nonstrict;
SET hive.txn.manager = org.apache.hadoop.hive.ql.lockmgr.DbTxnManager;
SET hive.compactor.initiator.on = true;
SET hive.compactor.worker.threads = 1;

创建Hive表如下:

create database logsdb;
use logsdb;

create table weblogs ( id int , msg string )

partitioned by (continent string, contry string)
clustered by (id) into 5 buckets
stored as orc
tblproperties(
  'transactional'='true',
  'transactional_properties'='default'
);

Flume Agnet 配置范例:

$ vi conf/hello-hive.conf 

#声明三种组件
a1.sources = r1
a1.channels = c1
a1.sinks = k1

#定义source信息
a1.sources.r1.type=netcat
a1.sources.r1.bind=localhost
a1.sources.r1.port=8888

#定义sink信息
a1.sinks.k1.type = hive
a1.sinks.k1.hive.metastore = thrift://hadoop100:9083
a1.sinks.k1.hive.database = logsdb
a1.sinks.k1.hive.table = weblogs
a1.sinks.k1.hive.partition = asia,India
a1.sinks.k1.useLocalTimeStamp = false
a1.sinks.k1.batchSize = 50
a1.sinks.k1.round = true
a1.sinks.k1.roundValue = 10
a1.sinks.k1.roundUnit = minute
a1.sinks.k1.serializer = DELIMITED
a1.sinks.k1.serializer.delimiter = "\t"
a1.sinks.k1.serializer.serdeSeparator = '\t'
a1.sinks.k1.serializer.fieldnames =id,msg

#定义channel信息
a1.channels.c1.type=memory

#绑定在一起
a1.sources.r1.channels=c1
a1.sinks.k1.channel=c1

启动Flume Agent

flume]$ flume-ng agent -c conf/ -f conf/hello-hive.conf -n a1 -Dflume.root.logger=INFO,console

使用nc客户端发送数据

[hadoop@hadoop100 ~]$ nc localhost 8888
1001    hello
OK
1002    world
OK

查看hive表中是否有数据出现

0: jdbc:hive2://localhost:10000> select * from weblogs;
+-------------+--------------+--------------------+-----------------+
| weblogs.id  | weblogs.msg  | weblogs.continent  | weblogs.contry  |
+-------------+--------------+--------------------+-----------------+
| 1002        | world        | asia               | India           |
| 1001        | hello        | asia               | India           |
+-------------+--------------+--------------------+-----------------+
2 rows selected (5.394 seconds)

如果对于行级更新删除需求比较频繁的,可以考虑使用事务表,但平常的hive表并不建议使用事务表。因为事务表的限制很多,加上由于hive表的特性,也很难满足高并发的场景。另外,如果事务表太多,并且存在大量的更新操作,metastore后台启动的合并线程会定期的提交MapReduce Job,也会一定程度上增重集群的负担。

Logger Sink

使用INFO级别把Event内容输出到日志中,一般用来测试、调试使用。这个 Sink 是唯一一个不需要额外配置就能把 Event 的原始内容输出的Sink,参照 输出原始数据到日志 。

提示

在 输出原始数据到日志 一节中说过,通常在Flume的运行日志里面输出数据流中的原始的数据内容是非常不可取的,所以 Flume 的组件默认都不会这么做。但是总有特殊的情况想要把 Event 内容打印出来,就可以借助这个Logger Sink了。

必需的参数已用 粗体 标明。

属性默认值解释
channel与 Sink 绑定的 channel
type组件类型,这个是: logger
maxBytesToLog16Event body 输出到日志的最大字节数,超出的部分会被丢弃

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = logger
a1.sinks.k1.channel = c1

Avro Sink

这个Sink可以作为 Flume 分层收集特性的下半部分。发送到此Sink的 Event 将转换为Avro Event发送到指定的主机/端口上。Event 从 channel 中批量获取,数量根据配置的 batch-size 而定。 必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channel
type 组件类型,这个是: avro.
hostname 监听的服务器名(hostname)或者 IP
port 监听的端口
batch-size 100 每次批量发送的 Event 数
connect-timeout 20000 第一次连接请求(握手)的超时时间,单位:毫秒
request-timeout 20000 请求超时时间,单位:毫秒
reset-connection-interval none 重置连接到下一跳之前的时间量(秒)。 这将强制 Avro Sink 重新连接到下一跳。 这将允许Sink在添加了新的主机时连接到硬件负载均衡器后面的主机,而无需重新启动 Agent。
compression-type none 压缩类型。可选值: none 、 deflate 。压缩类型必须与上一级Avro Source 配置的一致
compression-level 6 Event的压缩级别 0:不压缩,1-9:进行压缩,数字越大,压缩率越高
ssl false 设置为 true 表示开启SSL 下面的 truststore 、 truststore-password 、 truststore-type 就是开启SSL后使用的参数,并且可以指定是否信任所有证书( trust-all-certs )
trust-all-certs false 如果设置为true, 不会检查远程服务器(Avro Source)的SSL服务器证书。不要在生产环境开启这个配置,因为它使攻击者更容易执行中间人攻击并在加密的连接上进行“监听”。
truststore 自定义 Java truststore文件的路径。 Flume 使用此文件中的证书颁发机构信息来确定是否应该信任远程 Avro Source 的 SSL 身份验证凭据。 如果未指定,将使用全局的keystore配置,如果全局的keystore也未指定,将使用缺省 Java JSSE 证书颁发机构文件(通常为 Oracle JRE 中的“jssecacerts”或“cacerts”)。
truststore-password 上面配置的truststore的密码,如果未配置,将使用全局的truststore配置(如果配置了的话)
truststore-type JKS Java truststore的类型。可以配成 JKS 或者其他支持的 Java truststore类型,如果未配置,将使用全局的SSL配置(如果配置了的话)
exclude-protocols SSLv3 要排除的以空格分隔的 SSL/TLS 协议列表。 SSLv3 协议不管是否配置都会被排除掉。
maxIoWorkers 2 * 机器上可用的处理器核心数量 I/O工作线程的最大数量。这个是在 NettyAvroRpcClient 的 NioClientSocketChannelFactory 上配置的。

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = avro
a1.sinks.k1.channel = c1
a1.sinks.k1.hostname = 10.10.10.10
a1.sinks.k1.port = 4545

利用AvroSource和AvroSink实现跃点Agent

配置范例 - Agent #a1:

#a1
a1.sources = r1
a1.sinks= k1
a1.channels = c1

a1.sources.r1.type=netcat
a1.sources.r1.bind=localhost
a1.sources.r1.port=8888

a1.sinks.k1.type = avro
a1.sinks.k1.hostname=localhost
a1.sinks.k1.port=9999

a1.channels.c1.type=memory

a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

配置范例 - Agent #a2:

#a2
a2.sources = r2
a2.sinks= k2
a2.channels = c2

a2.sources.r2.type=avro
a2.sources.r2.bind=localhost
a2.sources.r2.port=9999

a2.sinks.k2.type = logger

a2.channels.c2.type=memory

a2.sources.r2.channels = c2
a2.sinks.k2.channel = c2

启动Agent #a2

$ flume-ng agent -f /soft/flume/conf/avro_hop.conf -n a2 -Dflume.root.logger=INFO,console

验证Agent #a2

$ netstat -anop | grep 9999

启动Agent #a1

$ flume-ng agent -f /soft/flume/conf/avro_hop.conf -n a1

验证Agent #a1

$ netstat -anop | grep 8888

HBase2Sink

提示

这是Flume 1.9新增的Sink。

HBase2Sink 是HBaseSink的HBase 2版本。

所提供的功能和配置参数与HBaseSink相同

必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channel
type 组件类型,这个是: hbase2
table 要写入的 Hbase 表名
columnFamily 要写入的 Hbase 列族
zookeeperQuorum Zookeeper 节点(host:port格式,多个用逗号分隔),hbase-site.xml 中属性 hbase.zookeeper.quorum 的值
znodeParent /hbase ZooKeeper 中 HBase 的 Root ZNode 路径,hbase-site.xml 中 zookeeper.znode.parent 的值
batchSize 100 每个事务写入的Event数量
coalesceIncrements false 每次提交时,Sink是否合并多个 increment 到一个cell。如果有限数量的 cell 有多个 increment ,这样可能会提供更好的性能
serializer org.apache.flume.sink.hbase2.SimpleHBase2EventSerializer 默认的列 increment column = “iCol”, payload column = “pCol”
serializer.* 序列化器的一些属性
kerberosPrincipal 以安全方式访问 HBase 的 Kerberos 用户主体
kerberosKeytab 以安全方式访问 HBase 的 Kerberos keytab 文件目录

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = hbase2
a1.sinks.k1.table = foo_table
a1.sinks.k1.columnFamily = bar_cf
a1.sinks.k1.serializer = org.apache.flume.sink.hbase2.RegexHBase2EventSerializer
a1.sinks.k1.channel = c1

Kafka Sink

这个 Sink 可以把数据发送到 Kafka topic上。目的就是将 Flume 与 Kafka 集成,以便基于拉的处理系统可以处理来自各种 Flume Source 的数据。

目前支持Kafka 0.10.1.0以上版本,最高已经在Kafka 2.0.1版本上完成了测试,这已经是Flume 1.9发行时候的最高的Kafka版本了。

必需的参数已用 粗体 标明。

属性 默认值 解释
type 组件类型,这个是: org.apache.flume.sink.kafka.KafkaSink
kafka.bootstrap.servers Kafka Sink 使用的 Kafka 集群的实例列表,可以是实例的部分列表。但是更建议至少两个用于高可用(HA)支持。格式为 hostname:port,多个用逗号分隔
kafka.topic default-flume-topic 用于发布消息的 Kafka topic 名称 。如果这个参数配置了值,消息就会被发布到这个 topic 上。如果Event header中包含叫做“topic”的属性, Event 就会被发布到 header 中指定的 topic 上,而不会发布到 kafka.topic 指定的 topic 上。支持任意的 header 属性动态替换, 比如%{lyf}就会被 Event header 中叫做“lyf”的属性值替换(如果使用了这种动态替换,建议将 Kafka 的 auto.create.topics.enable 属性设置为 true )。
flumeBatchSize 100 一批中要处理的消息数。设置较大的值可以提高吞吐量,但是会增加延迟。
kafka.producer.acks 1 在考虑成功写入之前,要有多少个副本必须确认消息。可选值, 0 :(从不等待确认); 1 :只等待leader确认; -1 :等待所有副本确认。 设置为-1可以避免某些情况 leader 实例失败的情况下丢失数据。
useFlumeEventFormat false 默认情况下,会直接将 Event body 的字节数组作为消息内容直接发送到 Kafka topic 。如果设置为true,会以 Flume Avro 二进制格式进行读取。 与 Kafka Source 上的同名参数或者 Kafka channel 的 parseAsFlumeEvent 参数相关联,这样以对象的形式处理能使生成端发送过来的 Event header 信息得以保留。
defaultPartitionId 指定所有 Event 将要发送到的 Kafka 分区ID,除非被 partitionIdHeader 参数的配置覆盖。 默认情况下,如果没有设置此参数,Event 会被 Kafka 生产者的分发程序分发,包括 key(如果指定了的话),或者被 kafka.partitioner.class 指定的分发程序来分发
partitionIdHeader 设置后,Sink将使用 Event header 中使用此属性的值命名的字段的值,并将消息发送到 topic 的指定分区。 如果该值表示无效分区,则将抛出 EventDeliveryException。 如果存在标头值,则此设置将覆盖 defaultPartitionId 。假如这个参数设置为“lyf”,这个 Sink 就会读取 Event header 中的 lyf 属性的值,用该值作为分区ID
allowTopicOverride true 如果设置为 true,会读取 Event header 中的名为 topicHeader 的的属性值,用它作为目标 topic。
topicHeader topic 与上面的 allowTopicOverride 一起使用,allowTopicOverride 会用当前参数配置的名字从 Event header 获取该属性的值,来作为目标 topic 名称
kafka.producer.security.protocol PLAINTEXT 设置使用哪种安全协议写入 Kafka。可选值:SASL_PLAINTEXT 、 SASL_SSL 和 SSL, 有关安全设置的其他信息,请参见下文。
more producer security props   如果使用了 SASL_PLAINTEXT 、 SASL_SSL 或 SSL 等安全协议,参考 Kafka security 来为生产者增加安全相关的参数配置
Other Kafka Producer Properties 其他一些 Kafka 生产者配置参数。任何 Kafka 支持的生产者参数都可以使用。唯一的要求是使用“kafka.producer.”这个前缀来配置参数,比如:kafka.producer.linger.ms

注解

Kafka Sink使用 Event header 中的 topic 和其他关键属性将 Event 发送到 Kafka。 如果 header 中存在 topic,则会将Event发送到该特定 topic,从而覆盖为Sink配置的 topic。 如果 header 中存在指定分区相关的参数,则Kafka将使用相关参数发送到指定分区。 header中特定参数相同的 Event 将被发送到同一分区。 如果为空,则将 Event 会被发送到随机分区。 Kafka Sink 还提供了key.deserializer(org.apache.kafka.common.serialization.StringSerializer) 和value.deserializer(org.apache.kafka.common.serialization.ByteArraySerializer)的默认值,不建议修改这些参数。

弃用的一些参数:

属性 默认值 解释
brokerList 改用 kafka.bootstrap.servers
topic default-flume-topic 改用 kafka.topic
batchSize 100 改用 kafka.flumeBatchSize
requiredAcks 1 改用 kafka.producer.acks

下面给出 Kafka Sink 的配置示例。Kafka 生产者的属性都是以 kafka.producer 为前缀。Kafka 生产者的属性不限于下面示例的几个。此外,可以在此处包含您的自定义属性,并通过作为方法参数传入的Flume Context对象在预处理器中访问它们。

a1.sinks.k1.channel = c1
a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink
a1.sinks.k1.kafka.topic = mytopic
a1.sinks.k1.kafka.bootstrap.servers = localhost:9092
a1.sinks.k1.kafka.flumeBatchSize = 20
a1.sinks.k1.kafka.producer.acks = 1
a1.sinks.k1.kafka.producer.linger.ms = 1
a1.sinks.k1.kafka.producer.compression.type = snappy

HTTP Sink

HTTP Sink 从 channel 中获取 Event,然后再向远程 HTTP 接口 POST 发送请求,Event 内容作为 POST 的正文发送。

错误处理取决于目标服务器返回的HTTP响应代码。 Sink的 退避 和 就绪 状态是可配置的,事务提交/回滚结果以及Event是否发送成功在内部指标计数器中也是可配置的。

状态代码不可读的服务器返回的任何格式错误的 HTTP 响应都将产生 退避 信号,并且不会从 channel 中消耗该Event。

必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channel
type 组件类型,这个是: http.
endpoint 将要 POST 提交数据接口的绝对地址
connectTimeout 5000 连接超时(毫秒)
requestTimeout 5000 一次请求操作的最大超时时间(毫秒)
contentTypeHeader text/plain HTTP请求的Content-Type请求头
acceptHeader text/plain HTTP请求的Accept 请求头
defaultBackoff true 是否默认启用退避机制,如果配置的 backoff.CODE 没有匹配到某个 http 状态码,默认就会使用这个参数值来决定是否退避
defaultRollback true 是否默认启用回滚机制,如果配置的 rollback.CODE 没有匹配到某个 http 状态码,默认会使用这个参数值来决定是否回滚
defaultIncrementMetrics false 是否默认进行统计计数,如果配置的 incrementMetrics.CODE 没有匹配到某个 http 状态码,默认会使用这个参数值来决定是否参与计数
backoff.CODE 配置某个 http 状态码是否启用退避机制(支持200这种精确匹配和2XX一组状态码匹配模式)
rollback.CODE 配置某个 http 状态码是否启用回滚机制(支持200这种精确匹配和2XX一组状态码匹配模式)
incrementMetrics.CODE 配置某个 http 状态码是否参与计数(支持200这种精确匹配和2XX一组状态码匹配模式)

注意 backoff, rollback 和 incrementMetrics 的 code 配置通常都是用具体的HTTP状态码,如果2xx和200这两种配置同时存在,则200的状态码会被精确匹配,其余200~299(除了200以外)之间的状态码会被2xx匹配。

提示

Flume里面好多组件都有这个退避机制,其实就是下一级目标没有按照预期执行的时候,会执行一个延迟操作。比如向HTTP接口提交数据发生了错误触发了退避机制生效,系统等待30秒再执行后续的提交操作, 如果再次发生错误则等待的时间会翻倍,直到达到系统设置的最大等待上限。通常在重试成功后退避就会被重置,下次遇到错误重新开始计算等待的时间。

任何空的或者为 null 的 Event 不会被提交到HTTP接口上。

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = http
a1.sinks.k1.channel = c1
a1.sinks.k1.endpoint = http://localhost:8080/someuri
a1.sinks.k1.connectTimeout = 2000
a1.sinks.k1.requestTimeout = 2000
a1.sinks.k1.acceptHeader = application/json
a1.sinks.k1.contentTypeHeader = application/json
a1.sinks.k1.defaultBackoff = true
a1.sinks.k1.defaultRollback = true
a1.sinks.k1.defaultIncrementMetrics = false
a1.sinks.k1.backoff.4XX = false
a1.sinks.k1.rollback.4XX = false
a1.sinks.k1.incrementMetrics.4XX = true
a1.sinks.k1.backoff.200 = false
a1.sinks.k1.rollback.200 = false
a1.sinks.k1.incrementMetrics.200 = true

File Roll Sink

把 Event 存储到本地文件系统。 必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channel
type 组件类型,这个是: file_roll.
sink.directory Event 将要保存的目录
sink.pathManager DEFAULT 配置使用哪个路径管理器,这个管理器的作用是按照规则生成新的存储文件名称,可选值有: default 、 rolltime。default规则:prefix+当前毫秒值+“-”+文件序号+“.”+extension;rolltime规则:prefix+yyyyMMddHHmmss+“-”+文件序号+“.”+extension;注:prefix 和 extension 如果没有配置则不会附带
sink.pathManager.extension 如果上面的 pathManager 使用默认的话,可以用这个属性配置存储文件的扩展名
sink.pathManager.prefix 如果上面的 pathManager 使用默认的话,可以用这个属性配置存储文件的文件名的固定前缀
sink.rollInterval 30 表示每隔30秒创建一个新文件进行存储。如果设置为0,表示所有 Event 都会写到一个文件中。
sink.serializer TEXT 配置 Event 序列化器,可选值有:text 、 header_and_text 、 avro_event 或者自定义实现了 EventSerializer.Builder 接口的序列化器的全限定类名.。 text 只会把 Event 的 body 的文本内容序列化; header_and_text 会把 header 和 body 内容都序列化。
sink.batchSize 100 每次事务批处理的 Event 数

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = file_roll
a1.sinks.k1.channel = c1
a1.sinks.k1.sink.directory = /var/log/flume

Null Sink

丢弃所有从 channel 读取到的 Event。 必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channel
type 组件类型,这个是: null.
batchSize 100 每次批处理的 Event 数量

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = null
a1.sinks.k1.channel = c1

Custom Sink

你可以自己写一个 Sink 接口的实现类。启动 Flume 时候必须把你自定义 Sink 所依赖的其他类配置进 classpath 内。custom source 在写配置文件的 type 时候填你的全限定类名。 必需的参数已用 粗体 标明。

属性 默认值 解释
channel 与 Sink 绑定的 channe
type 组件类型,这个填你自定义class的全限定类名

配置范例:

a1.channels = c1
a1.sinks = k1
a1.sinks.k1.type = org.example.MySink
a1.sinks.k1.channel = c1

Views: 156

HBase(十四)SQL引擎Phoenix

Phoenix介绍

1.什么是Phoenix

Phoenix是一个HBase的开源SQL引擎。你可以使用标准的JDBC API代替HBase客户端API来创建表,插入数据,查询你的HBase数据。

2.Phoenix底层原理

Phoenix框架将命令行上键入的sql语句翻译成hbase指令,然后hbase用翻译好的指令去操作集群,执行完之后给客户端反馈结果。

3.安装部署

  • 需要先安装好hbase集群,phoenix只是一个工具,只需要在一台机器上安装就可以了,这里我们选择node02服务器来进行安装一台即可

1、下载安装包

2、上传解压

  • 将安装包上传到node02服务器的/kkb/soft路径下,然后进行解压
cd /kkb/soft/
tar -zxf phoenix-hbase-2.2-5.1.1-bin.tar.gz   -C /kkb/install/

3.安装

cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin
cp -a phoenix-server-hbase-2.2-5.1.1.jar ../hbase-2.2.6/lib/
mv bin/hbase-site.xml  bin/hbase-site.xml.init
cp $HBASE_HOME/conf/hbase-site.xml ./bin
cp $HADOOP_HOME/etc/hadoop/hdfs-site.xml ./bin

4.配置环境变量

sudo vim /etc/profile
#phoenix
export PHOENIX_HOME=/kkb/install/phoenix-hbase-2.2-5.1.1-bin
export PHOENIX_CLASSPATH=$PHOENIX_HOME
export PATH=$PATH:$PHOENIX_HOME/bin
source /etc/profile

5.重启hbase集群

  • 记得要先启动hadoop集群、zookeeper集群

  • ==node01执行==以下命令来重启hbase的集群

cd /kkb/install/hbase-2.2.6
bin/start-hbase.sh
bin/stop-hbase.sh 

4、验证是否成功

  • 1、在phoenix/bin下输入命令, 进入到命令行,接下来就可以操作了

    ==node02执行==以下命令,进入phoenix客户端

    cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin/
    bin/sqlline.py node02:2181

5、Phoenix使用

1、 批处理方式

  • 1、node02执行以下命令创建user_phoenix.sql文件
    • 内容如下
mkdir -p /kkb/install/phoenixsql
cd /kkb/install/phoenixsql/
vim user_phoenix.sql

create table if not exists user_phoenix (state varchar(10) NOT NULL,  city varchar(20) NOT NULL, population BIGINT  CONSTRAINT my_pk PRIMARY KEY (state, city));
  • 2、node02执行以下命令,创建user_phoenix.csv数据文件
cd /kkb/install/phoenixsql/
vim user_phoenix.csv

NY,New York,8143197
CA,Los Angeles,3844829
IL,Chicago,2842518
TX,Houston,2016582
PA,Philadelphia,1463281
AZ,Phoenix,1461575
TX,San Antonio,1256509
CA,San Diego,1255540
TX,Dallas,1213825
CA,San Jose,912332
  • 3、创建user_phoenix_query.sql文件
cd /kkb/install/phoenixsql
vim user_phoenix_query.sql

select state as "userState",count(city) as "City Count",sum(population) as "Population Sum" FROM user_phoenix GROUP BY state; 
  • 4、执行sql语句
cd /kkb/install/phoenixsql

/kkb/install/phoenix-hbase-2.2-5.1.1-bin/bin/psql.py  node01:2181 user_phoenix.sql user_phoenix.csv  user_phoenix_query.sql

2、 命令行方式

  • 执行命令
cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin/
bin/sqlline.py node02:2181
  • 退出命令行方式,phoenix的命令都需要一个感叹号
!quit
  • 查看phoenix的帮助文档,显示所有命令
0: jdbc:phoenix:node02:2181> !help

2.1表的映射

  • 1、建立employee的映射表

    进入hbase客户端,创建一个普通表employee,并且有两个列族 company 和family
    node01执行以下以下命令进入hbase 的shell客户端

    cd /kkb/install/hbase-2.2.6/
    bin/hbase shell
    hbase(main):001:0> create 'employee','company','family'
  • 2、数据准备

    put 'employee','row1','company:name','ted'
    put 'employee','row1','company:position','worker'
    put 'employee','row1','family:tel','13600912345'
    put 'employee','row1','family:age','18'
    put 'employee','row2','company:name','michael'
    put 'employee','row2','company:position','manager'
    put 'employee','row2','family:tel','1894225698'
    put 'employee','row2','family:age','20'
  • 3、建立hbase到phoenix的映射表

    node02进入到phoenix的客户端,然后创建映射表

    cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin
    bin/sqlline.py node02:2181
    
    CREATE TABLE IF NOT EXISTS "employee" ("no" VARCHAR(10) NOT NULL PRIMARY KEY, "company"."name" VARCHAR(30),"company"."position" VARCHAR(20), "family"."tel" VARCHAR(20), "family"."age" VARCHAR(20)) column_encoded_bytes=0;

    说明

    在建立映射表之前要说明的是,Phoenix是大小写敏感的,并且所有命令都是大写

    如果你建的表名没有用双引号括起来,那么无论你输入的是大写还是小写,建立出来的表名都是大写的

    如果你需要建立出同时包含大写和小写的表名和字段名,请把表名或者字段名用双引号括起来

  • 4、查询映射表数据

0: jdbc:phoenix:node1:2181> select * from "employee";
+-------+----------+-----------+--------------+------+
|  no   |   name   | position  |     tel      | age  |
+-------+----------+-----------+--------------+------+
| row1  | ted      | worker    | 13600912345  | 18   |
| row2  | michael  | manager   | 1894225698   | 20   |
+-------+----------+-----------+--------------+------+

0: jdbc:phoenix:node01:2181> select * from "employee" where "tel" = '13600912345';
+-------+-------+-----------+--------------+------+
|  no   | name  | position  |     tel      | age  |
+-------+-------+-----------+--------------+------+
| row1  | ted   | worker    | 13600912345  | 18   |
+-------+-------+-----------+--------------+------+

3.JDBC

  • 创建maven工程并导入jar包
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <parent>
        <artifactId>XZK</artifactId>
        <groupId>org.example</groupId>
        <version>1.0-SNAPSHOT</version>
    </parent>
    <modelVersion>4.0.0</modelVersion>

    <artifactId>phoenixDemo</artifactId>
    <dependencies>
        <dependency>
            <groupId>org.apache.hbase</groupId>
            <artifactId>hbase-client</artifactId>
            <version>2.2.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.phoenix</groupId>
            <artifactId>phoenix-core</artifactId>
            <version>5.0.0-HBase-2.0</version>
        </dependency>
        <dependency>
            <groupId>junit</groupId>
            <artifactId>junit</artifactId>
            <version>4.12</version>
        </dependency>
        <dependency>
            <groupId>org.testng</groupId>
            <artifactId>testng</artifactId>
            <version>6.14.3</version>
        </dependency>
        <dependency>
            <groupId>junit</groupId>
            <artifactId>junit</artifactId>
            <version>4.13.1</version>
            <scope>test</scope>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <!-- 限制jdk版本插件 -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.0</version>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                    <encoding>UTF-8</encoding>
                </configuration>
            </plugin>
        </plugins>
    </build>

</project>
  • 代码开发
import org.junit.Before;
import org.junit.Test;

import java.sql.*;

public class PhoenixSearch {
    /**
     * 定义phoenix的url地址
     * connection
     * 构建Statement对象
     * 定义查询的sql语句,一定注意大小写
     * 构建好的对象执行sql语句
     */
    private Statement statement;
    private ResultSet rs;
    private Connection connection;

    @Before
    public void init() throws SQLException {
        //定义phoenix的url地址
        String url = "jdbc:phoenix:node02:2181";
        // connection
        connection = DriverManager.getConnection(url);
        //构建Statement对象
        statement = connection.createStatement();

    }

    @Test
    public void queryTable() throws SQLException {
        //定义查询的sql语句,一定注意大小写
        String sql="select * from USER_PHOENIX";
        //执行sql语句
        try {
            rs=statement.executeQuery(sql);
            while(rs.next()){
                System.out.println("city:"+rs.getString("city"));
                System.out.println("POPULATION :"+rs.getString("POPULATION"));
                System.out.println("STATE:"+rs.getString("STATE"));
                System.out.println("-------------------------");
            }
        } catch (SQLException e) {
            e.printStackTrace();
        } finally {
            if (connection != null) {
                connection.close();
            }
        }
    }
}

6.Phoenix构建二级索引

1、为什么需要用二级索引?

  • 对于HBase而言,如果想精确地定位到某行记录,唯一的办法是通过rowkey来查询。如果不通过rowkey来查找数据,就必须逐行地比较每一列的值,即全表扫瞄。

  • 对于较大的表,全表扫描的代价是不可接受的。但是,很多情况下,需要从多个角度查询数据。

    • 例如,在定位某个人的时候,可以通过姓名、身份证号、学籍号等不同的角度来查询
    • 要想把这么多角度的数据都放到rowkey中几乎不可能(业务的灵活性不允许,对rowkey长度的要求也不允许)。
    • 所以需要secondary index(二级索引)来完成这件事。secondary index的原理很简单,但是如果自己维护的话则会麻烦一些。
    • 现在,Phoenix已经提供了对HBase secondary index的支持。

2、Phoenix Global Indexing And Local Indexing

2.1 Global Indexing

  • Global indexing,全局索引,适用于读多写少的业务场景。
  • 使用Global indexing在写数据的时候开销很大,因为所有对数据表的更新操作(DELETE, UPSERT VALUES and UPSERT SELECT),都会引起索引表的更新,而索引表是分布在不同的数据节点上的,跨节点的数据传输带来了较大的性能消耗。
  • 在读数据的时候Phoenix会选择索引表来降低查询消耗的时间。
    • 在默认情况下如果想查询的字段不是索引字段的话索引表不会被使用,也就是说不会带来查询速度的提升。

2.2 Local Indexing

  • Local indexing,本地索引,适用于写操作频繁以及空间受限制的场景。
  • 与Global indexing一样,Phoenix会自动判定在进行查询的时候是否使用索引。
  • 使用Local indexing时,索引数据和数据表的数据存放在相同的服务器中,这样避免了在写操作的时候往不同服务器的索引表中写索引带来的额外开销。
  • 使用Local indexing的时候即使查询的字段不是索引字段索引表也会被使用,这会带来查询速度的提升,这点跟Global indexing不同。对于Local Indexing,一个数据表的所有索引数据都存储在一个单一的独立的可共享的表中。

3、Immutable index And Mutable index

3.1 immutable index

  • immutable index,不可变索引,适用于数据==只增加不更新并且按照时间先后顺序存储==(time-series data)的场景,如保存日志数据或者事件数据等。
  • 不可变索引的存储方式是write one,append only。
  • 当在Phoenix使用create table语句时指定IMMUTABLE_ROWS = true表示该表上创建的索引将被设置为不可变索引。
  • 不可变索引分为Global immutable index和Local immutable index两种。
  • Phoenix默认情况下如果在create table时不指定IMMUTABLE_ROW = true时,表示该表为mutable。

3.2 mutable index

  • mutable index,可变索引,适用于数据有增删改的场景。
  • Phoenix默认情况创建的索引都是可变索引,除非在create table的时候显式地指定IMMUTABLE_ROWS = true。
  • 可变索引同样分为Global mutable index和Local mutable index两种。

4、配置HBase支持Phoenix二级索引

4.1 修改配置文件

  • 如果要启用phoenix的二级索引功能,需要修改配置文件hbase-site.xml
  • 注意:
  • vim hbase-site.xml
<!-- 添加配置 -->
<property>
    <name>hbase.regionserver.wal.codec</name>
    <value>org.apache.hadoop.hbase.regionserver.wal.IndexedWALEditCodec</value>
</property>
<property>
   <name>hbase.region.server.rpc.scheduler.factory.class</name>
   <value>org.apache.hadoop.hbase.ipc.PhoenixRpcSchedulerFactory</value>
</property>
<property>
    <name>hbase.rpc.controllerfactory.class</name>
    <value>org.apache.hadoop.hbase.ipc.controller.ServerRpcControllerFactory</value>
</property>

4.2 重启hbase

  • 完成上述修改后重启hbase集群使配置生效。

5、实战

5.1 在phoenix中创建表

  • 首先,在phoenix中创建一个user table
  • node02执行以下命令,进入phoenix客户端,并创建表
cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin/
bin/sqlline.py node02:2181

create  table user (
"session_id" varchar(100) not null primary key, 
"f"."cookie_id" varchar(100), 
"f"."visit_time" varchar(100), 
"f"."user_id" varchar(100), 
"f"."age" varchar(100), 
"f"."sex" varchar(100), 
"f"."visit_url" varchar(100), 
"f"."visit_os" varchar(100), 
"f"."browser_name" varchar(100),
"f"."visit_ip" varchar(100), 
"f"."province" varchar(100),
"f"."city" varchar(100),
"f"."page_id" varchar(100), 
"f"."goods_id" varchar(100),
"f"."shop_id" varchar(100)) column_encoded_bytes=0;

5.2 导入测试数据

  • 将课件当中的user50w.csv 这个文件上传到node02的/kkb/install/phoenixsql 这个路径下
    该CSV文件中有50万条记录
  • node02执行以下命令,导入50W的测试数据
cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin/
bin/psql.py -t USER node01:2181 /kkb/install/phoenixsql/user50w.csv

5.3 Global Indexing的二级索引测试

5.3.1 正常查询一条数据所需的时间
  • 在为表USER创建secondary index之前,先看看查询一条数据所需的时间
    在node02服务器,进入phoenix的客户端,然后执行以下sql语句,查询数据,查看耗费时间
cd /kkb/install/phoenix-hbase-2.2-5.1.1-bin
bin/sqlline.py node01:2181

select * from user where "cookie_id" = '99738fd1-2084-44e9';

​ 可以看到,对名为cookie_id的列进行按值查询需要6秒左右。

​ 我们可以通过explain来查看执行计划
​ EXPLAIN(语句的执行逻辑及计划):

explain select * from user where "cookie_id" = '99738fd1-2084-44e9';

​ 由此看出先进行了全表扫描再通过过滤器来筛选出目标数据,显示这种查询方式效率是很低的。

5.3.2 给表USER创建基于Global Indexing的二级索引
  • 进入到phoenix的客户端,然后执行以下命令创建索引
  • 在cookie_id列上面创建二级索引:
0: jdbc:phoenix:node01:2181> create index USER_COOKIE_ID_INDEX on USER ("f"."cookie_id"); 

-- 查看当前所有表会发现多一张USER_COOKIE_ID_INDEX索引表,查询该表数据。
0: jdbc:phoenix:node01:2181> select * from USER_COOKIE_ID_INDEX limit 5;

  • 再次执行查询"cookie_id"='99738fd1-2084-44e9'的数据记录

    select "cookie_id" from user where "cookie_id" = '99738fd1-2084-44e9';

    此时:查询速度由10秒左右减少到了毫秒级别。

    注意:select所带的字段必须包含在覆盖索引内

    EXPLAIN(语句的执行逻辑及计划):

    explain select "cookie_id" from user where "cookie_id"='99738fd1-2084-44e9';

    可以看到使用到了创建的索引USER_COOKIE_ID_INDEX。

5.3.3 以下查询不会用到索引表
  • 虽然cookie_id是索引字段,但age不是索引字段,所以不会使用到索引
select "cookie_id","age" from user where "cookie_id"='99738fd1-2084-44e9';

  • 也可以通过EXPLAIN查询语句的执行逻辑及计划
EXPLAIN select "cookie_id","age" from user where "cookie_id"='99738fd1-2084-44e9';

  • 同理要查询的字段不是索引字段,也不会使用到索引表。
select "sex" from user where "cookie_id"='99738fd1-2084-44e9';

5.4 Local Indexing的二级索引测试

5.4.1 正常查询一条数据所需的时间
  • 在为表USER创建secondary index之前,先看看查询一条数据所需的时间
select * from user where "user_id"='371e963d-c-487065';
  • 可以看到,对名为user_id的列进行按值查询需要3秒左右。

  • EXPLAIN(语句的执行逻辑及计划):
explain select * from user where "user_id"='371e963d-c-487065';
  • 由此知道先进行了全表扫描再通过过滤器来筛选出目标数据,显示这种查询方式效率是很低的。

5.4.2 给表USER创建基于Local Indexing的二级索引
  • 在user_id列上面创建二级索引:
create local index USER_USER_ID_INDEX on USER ("f"."user_id");
  • 查看当前所有表会发现多一张USER_USER_ID_INDEX索引表,查询该表数据。

  • 再次执行查询"user_id"='371e963d-c-487065'的数据记录
select * from user where "user_id"='371e963d-c-487065';
  • 可以看到,对名为user_id的列进行按值查询需要0.7秒左右。

  • EXPLAIN(语句的执行逻辑及计划):
explain select * from user where "user_id"='371e963d-c-487065';

​ 查看执行计划,没有执行全表扫描,效率更高了

​ 此时:查询速度由3秒左右减少到了毫秒级别。

那如果查询的字段不包含在索引表中,又如何呢?

select "user_id","age","sex" from user where "user_id"='371e963d-c-487065';

EXPLAIN(语句的执行逻辑及计划):

explain select "user_id","age","sex" from user where "user_id"='371e963d-c-487065';

​ 可以看到使用到了创建的索引USER_USER_ID_INDEX.

5.5 如何确保query查询使用Index

  • 要想让一个查询使用index,有三种方式实现。
5.5.1 创建 covered index
  • 如果在某次查询中,查询项或者查询条件中包含除被索引列之外的列(主键MY_PK除外)。
  • 默认情况下,该查询会触发full table scan(全表扫描),但是使用covered index则可以避免全表扫描。创建包含某个字段的覆盖索引,创建方式如下:
create index USER_COOKIE_ID_AGE_INDEX on USER ("f"."cookie_id") include("f"."age");
  • 查看当前所有表会发现多一张USER_COOKIE_ID_AGE_INDEX索引表,查询该表数据。
select * from USER_COOKIE_ID_AGE_INDEX limit 5;

  • 查询数据
select "age" from user where "cookie_id"='99738fd1-2084-44e9';
select "age","sex" from user where "cookie_id"='99738fd1-2084-44e9';

5.5.2 在查询中提示其使用index
  • 在select和column_name之间加上/*+ Index(<表名> <index表名>)*/,通过这种方式强制使用索引。
select /*+ index(user,USER_COOKIE_ID_AGE_INDEX) */ "age" from user where "cookie_id"='99738fd1-2084-44e9';

如果age是索引字段,那么就会直接从索引表中查询
如果age不是索引字段,那么将会进行全表扫描,所以当用户明确知道表中数据较少且符合检索条件时才适用,此时的性能才是最佳的。

5.5.3 使用本地索引 (创建Local Indexing 索引)
  • 详细见上面

5.6 索引重建

  • Phoenix的索引重建是把索引表清空后重新装配数据。
alter index USER_COOKIE_ID_INDEX on user rebuild;

5.7 删除索引

  • 删除某个表的某张索引:
    语法 drop index 索引名称 on 表名
    例如:
drop  index USER_COOKIE_ID_INDEX on user;

​ 如果表中的一个索引列被删除,则索引也将被自动删除

​ 如果删除的是覆盖索引上的列,则此列将从覆盖索引中被自动删除。

6、索引性能调优

  • 一般来说,索引已经很快了,不需要特别的优化。
  • 这里也提供了一些方法,让你在面对特定的环境和负载的时候可以进行一些调优。下面的这些需要在hbase-site.xml文件中设置,针对所有的服务器。
1. index.builder.threads.max 
创建索引时,使用的最大线程数。 
默认值: 10。

2. index.builder.threads.keepalivetime 
创建索引的创建线程池中线程的存活时间,单位:秒。 
默认值: 60

3. index.writer.threads.max 
写索引表数据的写线程池的最大线程数。 
更新索引表可以用的最大线程数,也就是同时可以更新多少张索引表,数量最好和索引表的数量一致。 
默认值: 10

4. index.writer.threads.keepalivetime 
索引写线程池中,线程的存活时间,单位:秒。
默认值:60

5. hbase.htable.threads.max 
每一张索引表可用于写的线程数。 
默认值: 2,147,483,647

6. hbase.htable.threads.keepalivetime 
索引表线程池中线程的存活时间,单位:秒。 
默认值: 60

7. index.tablefactory.cache.size 
允许缓存的索引表的数量。 
增加此值,可以在写索引表时不用每次都去重复的创建htable,这个值越大,内存消耗越多。 
默认值: 10

8. org.apache.phoenix.regionserver.index.handler.count 
处理全局索引写请求时,可以使用的线程数。 
默认值: 30

Views: 35

HBase(十三)存储引擎设计与LSM-tree

前言

在数据存储的领域,有两大阵营,以B+tree为基础的关系型数据库,MySQL,SQLServer。以及以LSM-tree为基础的NoSQL key-value 存储, LevelDB。 LSM是(Log Structured Merge的简称)在分布式存储系统中通常会被设计成append-only的系统,LSM系统主要是顺序写优化,例如commit log等等,并作为分布式系统底层的基石。因此,要了解LSM的实现是十分重要的,下面主要介绍基于跳表的LSM-tree的实现 (Skiplist-Based LSM tree)。

为什么使用LSM-tree

讲LSM树之前,需要提下三种基本的存储引擎,这样才能清楚LSM树的由来:

  1. 哈希存储引擎 是哈希表的持久化实现,支持增、删、改以及随机读取操作,但不支持顺序扫描,对应的存储系统为key-value存储系统。对于key-value的插入以及查询,哈希表的复杂度都是O(1),明显比树的操作O(n)快,如果不需要有序的遍历数据,哈希表就是your Mr.Right

  2. B树存储引擎是B树(关于B树的由来,数据结构以及应用场景可以看之前一篇博文)的持久化实现,不仅支持单条记录的增、删、读、改操作,还支持顺序扫描(B+树的叶子节点之间的指针),对应的存储系统就是关系数据库(Mysql等)。

  3. LSM树(Log-Structured Merge Tree)存储引擎和B树存储引擎一样,同样支持增、删、读、改、顺序扫描操作。而且通过批量存储技术规避磁盘随机写入问题。当然凡事有利有弊,LSM树和B+树相比,LSM树牺牲了部分读性能,用来大幅提高写性能。

通过以上的分析,应该知道LSM树的由来了,LSM树的设计思想非常朴素:将对数据的修改增量保持在内存中,达到指定的大小限制后将这些修改操作批量写入磁盘。所以写入性能大大提升。

不过读取的时候需要合并磁盘中历史数据和内存中最近修改操作,读取时可能需要先看是否命中内存,否则需要访问较多的磁盘文件。

极端的说,基于LSM树实现的HBase的写性能比Mysql高了一个数量级,读性能低了一个数量级。

LSM树原理把一棵大树拆分成N棵小树,它首先写入内存中,随着小树越来越大,内存中的小树会flush到磁盘中,磁盘中的树定期可以做merge操作,合并成一棵大树,以优化读性能。

file

以上这些大概就是HBase存储的设计主要思想,这里分别对应说明下:

  1. 因为小树先写到内存中,为了防止内存数据丢失,写内存的同时需要暂时持久化到磁盘,对应了HBase的MemStore和HLog

  2. MemStore上的树达到一定大小之后,需要flush到HRegion磁盘中(一般是Hadoop DataNode),这样MemStore就变成DataNode上的磁盘文件StoreFile,定期HRegionServer对DataNode的数据做merge操作,彻底删除无效空间,多棵小树在这个时机合并成大树,来增强读性能。

LSM-Tree 原理

基于跳表的LSM-tree的实现 (Skiplist-Based LSM tree,以下简称为sLSM)

sLSM有两个宏观上(macrosocopic)的组件构成,内存中的buffer和磁盘的存储(in-memory buffer and disk-based store)。内存中的数据结构主要用于快速插入和对于buffered的数据进行查询。磁盘中的数据结构主要以层级的方式来存储数据,并提供merge的操作。

sLSM的实现,是基于Skiplist。即Memtable中使用Skiplist来快速查找。这里也可以有其他方式的实现,例如AVL, B-tree等。

LSM-tree在Level的实现如下图,可以看到主要有内存和磁盘的两大数据结构组件构成。内存中的数据结构支持缓存数据的查询(近期数据的查询),磁盘中的存储是层级形式的,且是不可变的(immuatble)。

file

关于LSM Tree,对于最简单的二层LSM Tree而言,内存中的数据和磁盘你中的数据merge操作,如下图

file

图来自lsm论文

lsm tree,理论上,可以是内存中树的一部分和磁盘中第一层树做merge,对于磁盘中的树直接做update操作有可能会破坏物理block的连续性,但是实际应用中,一般lsm有多层,当磁盘中的小树合并成一个大树的时候,可以重新排好顺序,使得block连续,优化读性能。

在hbase实现中,内存部分采用跳跃表来维护一个有序的KeyValue集合memstore,数据首先写入到memstore中,为了避免memstore内存数据丢失这里引入了WAL机制,会先将数据写入到所属regionServer的WAL日志中,也即HLog。memstore的数据达到某个级别的阈值之后,会进行写盘操作,形成一个有序的数据文件是(文件内部KeyValue有序),存储在磁盘上,磁盘部分则是对应的HFile。HFile就是一个小的树,因为hbase一般是部署在hdfs上,hdfs不支持对文件的update操作,而且最终随着磁盘文件越来越多,对读的影响很大。所以内存flush到磁盘上的小树,定期也会合并成一个大树。来增强读操作的性能,整体上hbase就是用了lsm tree的思路,如下图所示。

file

附加资料 - 跳表数据结构

SkipList(跳表) - 这种数据结构是由William Pugh于1990年在在 Communications of the ACM June 1990, 33(6) 668-676 发表了Skip lists: a probabilistic alternative to balanced trees,在其中详细描述了他的工作。由论文标题可知,SkipList的设计初衷是作为替换平衡树的一种选择。

什么是跳表

跳表全称为跳跃列表,它允许快速查询,插入和删除一个有序连续元素的数据链表。跳跃列表的平均查找和插入时间复杂度都是O(logn)。

快速查询是通过维护一个多层次的链表,且每一层链表中的元素是前一层链表元素的子集(见右边的示意图)。一开始时,算法在最稀疏的层次进行搜索,直至需要查找的元素在该层两个相邻的元素中间。这时,算法将跳转到下一个层次,重复刚才的搜索,直到找到需要查找的元素为止。

file

一张跳跃列表的示意图。每个带有箭头的框表示一个指针, 而每行是一个稀疏子序列的链表;底部的编号框(黄色)表示有序的数据序列。查找从顶部最稀疏的子序列向下进行, 直至需要查找的元素在该层两个相邻的元素中间。

跳表的演化过程

对于单链表来说,即使数据是已经排好序的,想要查询其中的一个数据,只能从头开始遍历链表,这样效率很低,时间复杂度很高,是 O(n)。

那我们有没有什么办法来提高查询的效率呢?我们可以为链表建立一个“索引”,这样查找起来就会更快,如下图所示,我们在原始链表的基础上,每两个结点提取一个结点建立索引,我们把抽取出来的结点叫做索引层或者索引,down 表示指向原始链表结点的指针。

file

现在如果我们想查找一个数据,比如说 15,我们首先在索引层遍历,当我们遍历到索引层中值为 14 的结点时,我们发现下一个结点的值为 17,所以我们要找的 15 肯定在这两个结点之间。这时我们就通过 14 结点的 down 指针,回到原始链表,然后继续遍历,这个时候我们只需要再遍历两个结点,就能找到我们想要的数据。好我们从头看一下,整个过程我们一共遍历了 7 个结点就找到我们想要的值,如果没有建立索引层,而是用原始链表的话,我们需要遍历 10 个结点。

通过这个例子我们可以看出来,通过建立一个索引层,我们查找一个基点需要遍历的次数变少了,也就是查询的效率提高了。

那么如果我们给索引层再加一层索引呢?遍历的结点会不会更少呢,效率会不会更高呢?我们试试就知道了。

file

现在我们再来查找 15,我们从第二级索引开始,最后找到 15,一共遍历了 6 个结点,果然效率更高。

当然,因为我们举的这个例子数据量很小,所以效率提升的不是特别明显,如果数据量非常大的时候,我们多建立几层索引,效率提升的将会非常的明显,感兴趣的可以自己试一下,这里我们就不举例子了。

这种通过对链表加多级索引的机构,就是跳表了。

跳表具体有多快

通过上边的例子我们知道,跳表的查询效率比链表高,那具体高多少呢?下面我们一起来看一下。

衡量一个算法的效率我们可以用时间复杂度,这里我们也用时间复杂度来比较一下链表和跳表。前面我们已经讲过了,链表的查询的时间复杂度为 O(n),那跳表的呢?

如果一个链表有 n 个结点,如果每两个结点抽取出一个结点建立索引的话,那么第一级索引的结点数大约就是 n/2,第二级索引的结点数大约为 n/4,以此类推第 m 级索引的节点数大约为 n/(2^m)。

假如一共有 m 级索引,第 m 级的结点数为两个,通过上边我们找到的规律,那么得出 n/(2^m)=2,从而求得 m=log(n)-1。如果加上原始链表,那么整个跳表的高度就是 log(n)。我们在查询跳表的时候,如果每一层都需要遍历 k 个结点,那么最终的时间复杂度就为 O(k*log(n))。

那这个 k 值为多少呢,按照我们每两个结点提取一个基点建立索引的情况,我们每一级最多需要遍历两个个结点,所以 k=2。为什么每一层最多遍历两个结点呢?

因为我们是每两个结点提取一个结点建立索引,最高一级索引只有两个结点,然后下一层索引比上一层索引两个结点之间增加了一个结点,也就是上一层索引两结点的中值,看到这里是不是想起来我们前边讲过的二分查找,每次我们只需要判断要找的值在不在当前结点和下一个结点之间即可。

file

如上图所示,我们要查询红色结点,我们查询的路线即黄线表示出的路径查询,每一级最多遍历两个结点即可。

所以跳表的查询任意数据的时间复杂度为 O(2*log(n)),前边的常数 2 可以忽略,为 O(log(n))。

跳表是用空间来换时间

跳表的效率比链表高了,但是跳表需要额外存储多级索引,所以需要的更多的内存空间。

跳表的空间复杂度分析并不难,如果一个链表有 n 个结点,如果每两个结点抽取出一个结点建立索引的话,那么第一级索引的结点数大约就是 n/2,第二级索引的结点数大约为 n/4,以此类推第 m 级索引的节点数大约为 n/(2^m),我们可以看出来这是一个等比数列。

这几级索引的结点总和就是 n/2+n/4+n/8…+8+4+2=n-2,所以跳表的空间复杂度为 o(n)。

那么我们有没有办法减少索引所占的内存空间呢?可以的,我们可以每三个结点抽取一个索引,或者没五个结点抽取一个索引。这样索引结点的数量减少了,所占的空间也就少了。

跳表的插入和删除

我们想要为跳表插入或者删除数据,我们首先需要找到插入或者删除的位置,然后执行插入或删除操作,前边我们已经知道了,跳表的查询的时间复杂度为 O(logn),因为找到位置之后插入和删除的时间复杂度很低,为 O(1),所以最终插入和删除的时间复杂度也为 O(longn)。

我么通过图看一下插入的过程。

file

删除操作的话,如果这个结点在索引中也有出现,我们除了要删除原始链表中的结点,还要删除索引中的。因为单链表中的删除操作需要拿到要删除结点的前驱结点,然后通过指针操作完成删除。所以在查找要删除的结点的时候,一定要获取前驱结点。当然,如果我们用的是双向链表,就不需要考虑这个问题了。

如果我们不停的向跳表中插入元素,就可能会造成两个索引点之间的结点过多的情况。结点过多的话,我们建立索引的优势也就没有了。所以我们需要维护索引与原始链表的大小平衡,也就是结点增多了,索引也相应增加,避免出现两个索引之间结点过多的情况,查找效率降低。

跳表是通过一个随机函数来维护这个平衡的,当我们向跳表中插入数据的的时候,我们可以选择同时把这个数据插入到索引里,那我们插入到哪一级的索引呢,这就需要随机函数,来决定我们插入到哪一级的索引中。

这样可以很有效的防止跳表退化,而造成效率变低。

跳表的代码实现

最后我们来看一下跳变用代码怎么实现。

package skiplist;

import java.util.Random;

/**
 * 跳表的一种实现方法。
 * 跳表中存储的是正整数,并且存储的是不重复的。
 */
public class SkipList {

  private static final int MAX_LEVEL = 16;

  private static final float SKIPLIST_P = 0.5f;

  private int levelCount = 1;

  private Node head = new Node();  // 带头链表

  private Random r = new Random();

  public Node find(int value) {
    Node p = head;
    for (int i = levelCount - 1; i >= 0; --i) {
      while (p.forwards[i] != null && p.forwards[i].data < value) {
        p = p.forwards[i];
      }
    }

    if (p.forwards[0] != null && p.forwards[0].data == value) {
      return p.forwards[0];
    } else {
      return null;
    }
  }

  public void insert(int value) {
    int level = randomLevel();
    Node newNode = new Node();
    newNode.data = value;
    newNode.maxLevel = level;
    Node update[] = new Node[level];
    for (int i = 0; i < level; ++i) {
      update[i] = head;
    }

    // record every level largest value which smaller than insert value in update[]
    Node p = head;
    for (int i = level - 1; i >= 0; --i) {
      while (p.forwards[i] != null && p.forwards[i].data < value) {
        p = p.forwards[i];
      }
      update[i] = p;// use update save node in search path
    }

    // in search path node next node become new node forwords(next)
    for (int i = 0; i < level; ++i) {
      newNode.forwards[i] = update[i].forwards[i];
      update[i].forwards[i] = newNode;
    }

    // update node hight
    if (levelCount < level) levelCount = level;
  }

  public void delete(int value) {
    Node[] update = new Node[levelCount];
    Node p = head;
    for (int i = levelCount - 1; i >= 0; --i) {
      while (p.forwards[i] != null && p.forwards[i].data < value) {
        p = p.forwards[i];
      }
      update[i] = p;
    }

    if (p.forwards[0] != null && p.forwards[0].data == value) {
      for (int i = levelCount - 1; i >= 0; --i) {
        if (update[i].forwards[i] != null && update[i].forwards[i].data == value) {
          update[i].forwards[i] = update[i].forwards[i].forwards[i];
        }
      }
    }
  }

 // 理论来讲,一级索引中元素个数应该占原始数据的 50%,二级索引中元素个数占 25%,三级索引12.5% ,一直到最顶层。
  // 因为这里每一层的晋升概率是 50%。对于每一个新插入的节点,都需要调用 randomLevel 生成一个合理的层数。
  // 该 randomLevel 方法会随机生成 1~MAX_LEVEL 之间的数,且 :
  //        50%的概率返回 1
  //        25%的概率返回 2
  //      12.5%的概率返回 3 ...
  private int randomLevel() {
    int level = 1;

    while (Math.random() < SKIPLIST_P && level < MAX_LEVEL)
      level += 1;
    return level;
  }

  public void printAll() {
    Node p = head;
    while (p.forwards[0] != null) {
      System.out.print(p.forwards[0] + " ");
      p = p.forwards[0];
    }
    System.out.println();
  }

  public class Node {
    private int data = -1;
    private Node forwards[] = new Node[MAX_LEVEL];
    private int maxLevel = 0;

    @Override
    public String toString() {
      StringBuilder builder = new StringBuilder();
      builder.append("{ data: ");
      builder.append(data);
      builder.append("; levels: ");
      builder.append(maxLevel);
      builder.append(" }");

      return builder.toString();
    }
  }

}

Views: 26