分布式存储引擎 – KYLIN

Apache Kylin 是一个开源的分布式存储引擎,最初由 eBay 开发贡献至开源 社区。它提供 Hadoop 之上的 SQL 查询接口及多维分析(OLAP)能力以支持大规 模数据,能够处理 TB 乃至 PB 级别的分析任务,能够在亚秒级查询巨大的 Hive 表,并支持高并发。

file

1.1、为什么要使用kylin

自从 10 年前 Hadoop 诞生以来,大数据的存储和批处理问题均得到了妥善解 决,而如何高速地分析数据也就成为了下一个挑战。于是各式各样的“SQL on Hadoop”技术应运而生,其中以 Hive 为代表,Impala、Presto、Phoenix、Drill、 SparkSQL 等紧随其后。它们的主要技术是“大规模并行处理”(Massive Parallel Processing,MPP)和“列式存储”(Columnar Storage)。

大规模并行处理可以调动多台机器一起进行并行计算,用线性增加的资源来 换取计算时间的线性下降。列式存储则将记录按列存放,这样做不仅可以在访问 时只读取需要的列,还可以利用存储设备擅长连续读取的特点,大大提高读取的 速率。这两项关键技术使得 Hadoop 上的 SQL 查询速度从小时提高到了分钟。 然而分钟级别的查询响应仍然离交互式分析的现实需求还很远。分析师敲入 查询指令,按下回车,还需要去倒杯咖啡,静静地等待查询结果。得到结果之后 才能根据情况调整查询,再做下一轮分析。如此反复,一个具体的场景分析常常 需要几小时甚至几天才能完成,效率低下。 这是因为大规模并行处理和列式存储虽然提高了计算和存储的速度,但并没 有改变查询问题本身的时间复杂度,也没有改变查询时间与数据量成线性增长的 关系这一事实。假设查询 1 亿条记录耗时 1 分钟,那么查询 10 亿条记录就需 10分钟,100 亿条记录就至少需要 1 小时 40 分钟。 当然,可以用很多的优化技术缩短查询的时间,比如更快的存储、更高效的压缩算法,等等,但总体来说,查询性能与数据量呈线性相关这一点是无法改变 的。虽然大规模并行处理允许十倍或百倍地扩张计算集群,以期望保持分钟级别 的查询速度,但购买和部署十倍或百倍的计算集群又怎能轻易做到,更何况还有 高昂的硬件运维成本。 另外,对于分析师来说,完备的、经过验证的数据模型比分析性能更加重要, 直接访问纷繁复杂的原始数据并进行相关分析其实并不是很友好的体验,特别是 在超大规模的数据集上,分析师将更多的精力花在了等待查询结果上,而不是在 更加重要的建立领域模型上。

1.2、kylin的使用场景

(1) 假如你的数据存储于 Hadoop 的 HDFS 分布式文件系统中,并且使用 Hive 来基于 HDFS 构建数据仓库系统,并进行数据分析,但是数据量巨大, 比如 PB 级别;

(2) 同时也使用 HBase 来进行数据的存储和利于 HBase 的行键实现数据 的快速查询;

(3) 数据分析平台的数据量逐日累积增加;

(4) 对于数据分析的维度大概 10 个左右。 如果类似于上述的场景,那么非常适合使用 Apache Kylin 来做大数据的多维分析。

1.3、kylin如何解决海量数据的查询问题

Apache Kylin 的初衷就是要解决千亿条、万亿条记录的秒级查询问 题,其中的关键就是要打破查询时间随着数据量成线性增长的这个规律。仔细思 考大数据 OLAP,可以注意到两个事实。

大数据查询要的一般是统计结果,是多条记录经过聚合函数计算后的统计 值。原始的记录则不是必需的,或者访问频率和概率都极低。

聚合是按维度进行的,由于业务范围和分析需求是有限的,有意义的维度 聚合组合也是相对有限的,一般不会随着数据的膨胀而增长。

基于以上两点,我们可以得到一个新的思路——“预计算”。应尽量多地预 先计算聚合结果,在查询时刻应尽量使用预算的结果得出查询结果,从而避免直 接扫描可能无限增长的原始记录。

举例来说,使用如下的 SQL 来查询 10 月 1 日那天销量最高的商品:

file

用传统的方法时需要扫描所有的记录,再找到 10 月 1 日的销售记录,然后

按商品聚合销售额,最后排序返回。假如 10 月 1 日有 1 亿条交易,那么查询必

须读取并累计至少 1 亿条记录,且这个查询速度会随将来销量的增加而逐步下降。如果日交易量提高一倍到 2 亿,那么查询执行的时间可能也会增加一倍。 而使用 预 计 算 的 方 法 则 会 事 先 按 维 度 [sell_date , item] 计 算  sum

(sell_amount)并存储下来,在查询时找到  10 月 1 日的销售商品就可以直接

排序返回了。读取的记录数最大不会超过维度[sell_date,item]的组合数。显 然这个数字将远远小于实际的销售记录,比如 10 月 1 日的 1 亿条交易包含了 100

万条商品,那么预计算后就只有 100 万条记录了,是原来的百分之一。并且这些 记录已经是按商品聚合的结果,因此又省去了运行时的聚合运算。从未来的发展 来看,查询速度只会随日期和商品数目的增长而变化,与销售记录的总数不再有 直接联系。假如日交易量提高一倍到 2 亿,但只要商品的总数不变,那么预计算 的结果记录总数就不会变,查询的速度也不会变。

“预计算”就是 Kylin 在“大规模并行处理”和“列式存储”之外,提供给大数据分析的第三个关键技术。

2、Kylin前置基础知识了解

2.1、数据仓库、OLAP 与 BI

数据仓库

数据仓库,英文名称 Data Warehouse,简称 DW。《数据仓库》一书中的定义 为:数据仓库就是面向主题的、集成的、相对稳定的、随时间不断变化(不同时 间)的数据集合,用以支持经营管理中的决策制定过程、数据仓库中的数据面向 主题,与传统数据库面向应用相对应。

利用数据仓库的方式存放的资料,具有一旦存入,便不会随时间发生变动的 特性,此外,存入的资料必定包含时间属性,通常一个数据仓库中会含有大量的 历史性资料,并且它可利用特定的分析方式,从其中发掘出特定的资讯。

OLAP

1、OLAP的基本概念

OLAP(Online Analytical Process),联机分析处理,以多维度的方式分 析数据,而且能够弹性地提供上卷(Roll-up)、下钻(Drill-down)和切片(Slice) 等操作,它是呈现集成性决策信息的方法,多用于决策支持系统、商务智能或数 据仓库。其主要的功能在于方便大规模数据分析及统计计算,可对决策提供参考 和支持。与之相区别的是联机交易处理(OLTP),联机交易处理,更侧重于基本 的、日常的事务处理,包括数据的增删改查。

OLAP 需要以大量历史数据为基础,再配合上时间点的差异,对多维 度及汇整型的信息进行复杂的分析。

OLAP 需要用户有主观的信息需求定义,因此系统效率较佳。

OLAP 的概念,在实际应用中存在广义和狭义两种不同的理解方式。广义上 的理解与字面上的意思相同,泛指一切不会对数据进行更新的分析处理。但更多 的情况下 OLAP 被理解为其狭义上的含义,即与多维分析相关,基于立方体(Cube) 计算而进行的分析。

OLAP(online analytical processing)是一种软件技术,它使分析人员能够迅速、一致、交互地从各个方面观察信息,以达到深入理解数据的目的。从各方面观察信息,也就是从不同的维度分析数据,因此OLAP也成为多维分析。

file

2、OLAP的类型

也可以分为ROLAP和MOLAP

file

3、OLAP  CUBE

file

4、CUBE与 Cuboid

file

BI

BI(Business Intelligence),即商务智能,指用现代数据仓库技术、在线 分析技术、数据挖掘和数据展现技术进行数据分析以实现商业价值。

2.2、事实表与维度表

事实表(Fact Table)是指存储有事实记录的表,如系统日志、销售记录等; 事实表的记录在不断地动态增长,所以它的体积通常远大于其他表。

维度表(Dimension Table)或维表,有时也称查找表(Lookup Table),是 分析事实的一种角度,是与事实表相对应的一种表;它保存了维度的属性值,可 以跟事实表做关联;相当于将事实表上经常重复出现的属性抽取、规范出来用一 张表进行管理。常见的维度表有:日期表(存储与日期对应的周、月、季度等的 属性)、地点表(包含国家、省/州、城市等属性)等。使用维度表有诸多好处, 具体如下。

·缩小了事实表的大小。

·便于维度的管理和维护,增加、删除和修改维度的属性,不必对事实表的 大量记录进行改动。

·维度表可以为多个事实表重用,以减少重复工作。

2.3、维度与度量

维度是指审视数据的角度,它通常是数据记录的一个属性,例如时间、地点 等。

度量是基于数据所计算出来的考量值;它通常是一个数值,如总销售额、不 同的用户数等。 分析人员往往要结合若干个维度来审查度量值,以便在其中找到变化规律。 在一个 SQL 查询中,Group By 的属性通常就是维度,而所计算的值则是度量。 如下面的示例:

file

  在上面的这个查询中,part_dt 和 lstg_site_id 是维度,sum(price)和

count(distinct seller_id)是度量。

file

2.4、数据仓库建模常用手段方式

星型模型:

星形模型中有一张事实表,以及零个或多个维度表;事实表与维度表通过主 键外键相关联,维度表之间没有关联,就像很多星星围绕在一个恒星周围,故取 名为星形模型。

file

雪花模型:

若将星形模型中某些维度的表再做规范,抽取成更细的维度表,然后让维

度表之间也进行关联,那么这种模型称为雪花模型。

file

星座模式

星座模式是星型模式延伸而来,星型模式是基于一张事实表的,而星座模式是基于多张事实表的,而且共享维度信息。

前面介绍的两种维度建模方法都是多维表对应单事实表,但在很多时候维度空间内的事实表不止一个,而一个维表也可能被多个事实表用到。在业务发展后期,绝大部分维度建模都采用的是星座模式。

file

注意:Kylin 只支持星形模型的数据集

2.5、数据立方体

Cube(或 Data Cube),即数据立方体,是一种常用于数据分析与索引的技术;它可以对原始数据建立多维度索引。通过 Cube 对数据进行分析,可以大大 加快数据的查询效率。

Cuboid 在 Kylin 中特指在某一种维度组合下所计算的数据。 给定一个数据模型,我们可以对其上的所有维度进行组合。对于 N 个维度来

说,组合的所有可能性共有 2 的 N 次方种。对于每一种维度的组合,将度量做 聚合运算,然后将运算的结果保存为一个物化视图,称为 Cuboid。

所有维度组合的 Cuboid 作为一个整体,被称为 Cube。所以简单来说,一个 Cube 就是许多按维度聚合的物化视图的集合。

下面来列举一个具体的例子。假定有一个电商的销售数据集,其中维度包括 时间(Time)、商品(Item)、地点(Location)和供应商(Supplier),度量为销 售额(GMV)。那么所有维度的组合就有 2 的 4 次方 =16 种,比如一维度(1D) 的组合有[Time]、[Item]、[Location]、[Supplier]4 种;二维度(2D)的组合 有[Time,Item]、[Time,Location]、[Time、Supplier]、[Item,Location]、 [Item,Supplier]、[Location,Supplier]6 种;三维度(3D)的组合也有 4 种; 最后零维度(0D)和四维度(4D)的组合各有 1 种,总共就有 16 种组合。

file

2.6、Kylin的工作原理

Apache Kylin 的工作原理就是对数据模型做 Cube 预计算,并利用计算的结 果加速查询,具体工作过程如下。

1)指定数据模型,定义维度和度量。

2)预计算 Cube,计算所有 Cuboid 并保存为物化视图。

3)执行查询时,读取 Cuboid,运算,产生查询结果。

由于 Kylin 的查询过程不会扫描原始记录,而是通过预计算预先完成表的关 联、聚合等复杂运算,并利用预计算的结果来执行查询,因此相比非预计算的查 询技术,其速度一般要快一到两个数量级,并且这点在超大的数据集上优势更明 显。当数据集达到千亿乃至万亿级别时,Kylin 的速度甚至可以超越其他非预计算技术 1000 倍以上。

2.7、Kylin的体系架构

Apache Kylin 系统可以分为在线查询和离线构建两部分,技术架构如图所 示,在线查询的模块主要处于上半区,而离线构建则处于下半区。

file

1)REST Server

REST Server是一套面向应用程序开发的入口点,旨在实现针对Kylin平台的应用开发工作。 此类应用程序可以提供查询、获取结果、触发cube构建任务、获取元数据以及获取用户权限等等。另外可以通过Restful接口实现SQL查询。

2)查询引擎(Query Engine)

当cube准备就绪后,查询引擎就能够获取并解析用户查询。它随后会与系统中的其它组件进行交互,从而向用户返回对应的结果。 

3)路由器(Routing)

在最初设计时曾考虑过将Kylin不能执行的查询引导去Hive中继续执行,但在实践后发现Hive与Kylin的速度差异过大,导致用户无法对查询的速度有一致的期望,很可能大多数查询几秒内就返回结果了,而有些查询则要等几分钟到几十分钟,因此体验非常糟糕。最后这个路由功能在发行版中默认关闭。

4)元数据管理工具(Metadata)

Kylin是一款元数据驱动型应用程序。元数据管理工具是一大关键性组件,用于对保存在Kylin当中的所有元数据进行管理,其中包括最为重要的cube元数据。其它全部组件的正常运作都需以元数据管理工具为基础。 Kylin的元数据存储在hbase中。 

5)任务引擎(Cube Build Engine)

这套引擎的设计目的在于处理所有离线任务,其中包括shell脚本、Java API以及Map Reduce任务等等。任务引擎对Kylin当中的全部任务加以管理与协调,从而确保每一项任务都能得到切实执行并解决其间出现的故障。

2.8、Kylin特点

Kylin的主要特点包括支持SQL接口、支持超大规模数据集、亚秒级响应、可伸缩性、高吞吐率、BI工具集成等。

1)标准SQL接口:Kylin是以标准的SQL作为对外服务的接口。

2)支持超大数据集:Kylin对于大数据的支撑能力可能是目前所有技术中最为领先的。早在2015年eBay的生产环境中就能支持百亿记录的秒级查询,之后在移动的应用场景中又有了千亿记录秒级查询的案例。

3)亚秒级响应:Kylin拥有优异的查询相应速度,这点得益于预计算,很多复杂的计算,比如连接、聚合,在离线的预计算过程中就已经完成,这大大降低了查询时刻所需的计算量,提高了响应速度。

4)可伸缩性和高吞吐率:单节点Kylin可实现每秒70个查询,还可以搭建Kylin的集群。

5)BI工具集成

Kylin可以与现有的BI工具集成,具体包括如下内容。

ODBC:与Tableau、Excel、PowerBI等工具集成

JDBC:与Saiku、BIRT等Java工具集成

RestAPI:与JavaScript、Web网页集成

Kylin开发团队还贡献了Zepplin的插件,也可以使用Zepplin来访问Kylin服务。

3、Kylin的环境安装

1)官网地址

http://kylin.apache.org/cn/

2)官方文档

http://kylin.apache.org/cn/docs/

3)下载地址

http://kylin.apache.org/cn/download/

3.1 单节点服务模式安装

kylin的运行环境分为单机模式和集群模式,单机模式只需要在任意一台机器安装一台kylin服务即可,集群模式可以在所有机器上面都安装,然后所有机器的kylin组成集群

kylin的服务安装需要依赖于 zookeeper,hdfs,yarn,hive,hbase等各种服务,在安装kylin之前需要保证我们的zookeeper,hdfs,yarn,hive以及hbase的服务都是正常的并且是处于运行状态。

主机名 服务Node01Node02Node03
zookeeperQuorumPeerMainQuorumPeerMainQuorumPeerMain
hdfsnamenode  
secondaryNameNode  
DataNodeDataNodeDataNode
YarnResourceManager  
NodeManagerNodeManagerNodeManager
MapReduceJobHistoryServer  
HBaseHMaster  
HRegionServerHRegionServerHRegionServer
Hive  HiveServer2
  MetaStore

第一步:下载kylin安装包上传并解压

kylin安装包下载地址为

http://mirrors.tuna.tsinghua.edu.cn/apache/kylin/apache-kylin-2.6.3/apache-kylin-2.6.3-bin-cdh57.tar.gz

(如果你使用的是hadoop3版本,则推荐使用kylin3:https://mirrors.tuna.tsinghua.edu.cn/apache/kylin/apache-kylin-3.1.2/apache-kylin-3.1.2-bin-hadoop3.tar.gz)

将安装包上传到node03服务器的/kkb/soft路径下,并解压到/kkb/install

node03执行以下命令,进行解压

cd /kkb/soft
tar -zxf apache-kylin-2.6.3-bin-cdh57.tar.gz  -C /kkb/install/

第二步:node03服务器开发环境变量配置

node03服务器添加以下环境变量:

sudo vim /etc/profile

export JAVA_HOME=/kkb/install/jdk1.8.0_141
export PATH=:$JAVA_HOME/bin:$PATH

export HADOOP_HOME=/kkb/install/hadoop-2.6.0-cdh5.14.2
export PATH=:$HADOOP_HOME/bin:$PATH

export HBASE_HOME=/kkb/install/hbase-1.2.0-cdh5.14.2
export PATH=:$HBASE_HOME/bin:$PATH

export HIVE_HOME=/kkb/install/hive-1.1.0-cdh5.14.2
export PATH=:$HIVE_HOME/bin:$PATH

export HCAT_HOME=/kkb/install/hive-1.1.0-cdh5.14.2
export PATH=:$HCAT_HOME/hcatalog:$PATH

export KYLIN_HOME=/kkb/install/apache-kylin-2.6.3-bin-cdh57
export PATH=:$KYLIN_HOME/bin:$PATH

export dir=/kkb/install/apache-kylin-2.6.3-bin-cdh57/bin
export PATH=$dir:$PATH

更改完了环境变量,记得source /etc/profile 生效

第三步:node03启动kylin服务

node03执行以下命令启动kylin服务

cd /kkb/install/apache-kylin-2.6.3-bin-cdh57
bin/kylin.sh start

Kylin启动报错hbase-common lib not found

解决方案:修改hbase文件$HBASE_HOME/bin/hbase

找到下面这行

CLASSPATH=${CLASSPATH}:$JAVA_HOME/lib/tools.jar

修改为

CLASSPATH=${CLASSPATH}:$JAVA_HOME/lib/tools.jar:/YOUR HBASE FULL PATH or $HBASE_HOME/lib/*

第四步:浏览器访问kylin服务

浏览器界面访问kylin服务     http://node03.kaikeba.com:7070/kylin/

用户名:ADMIN   密码:KYLIN

3.2 kylin的集群环境安装

单节点的kylin环境,主要用于我们方便测试学习,实际工作当中,我们主要还是使用kylin的集群模式来进行开发,接下来我们就来看一下kylin的集群模式该如何运行

Kylin的实例是无状态的,运行时的状态保存在Hbase的元数据中(kylin.metadata.url指定)

只要每个实例都指向读取共同的元数据就可以完成集群的部署(即元数据共享)

对于每个实例,都必须指定实例运行的模式(kylin.server.mode),共有3种模式

job 只能运行job引擎

query 只能运行查询引擎

all 既可以运行job 又可以运行query

query模式下只支持sql查询,不执行cube的构建等相关操作。 特别注意:kylin集群中只能有一个实例运行job引擎,其他必须是query模式。

file

集群模式重要配置参数介绍

当kylin以集群模式运行的时候,会存在多个运行实例,可以通过conf/kylin.properties中两个参数进行设置

kylin.server.cluster-servers

列出所有rest   web  Servers,使得实例之间进行同步,比如设置为:

kylin.server.cluster-servers=node01:7070,node02:7070,node03:7070
kylin.server.mode

确保一个实例配置的是all或者job,其他都必须是query模式。

第一步:将node03服务器的kylin安装包分发到其他机器

将node03服务器/kkb/install路径下的kylin的安装包分发到其他服务器上面去

node03执行以下命令停止kylin服务,然后将kylin安装包分发到其他服务器上面去

node03执行以下命令

cd /kkb/install/apache-kylin-2.6.3-bin-cdh57
bin/kylin.sh stop

cd /kkb/install/
scp -r apache-kylin-2.6.3-bin-cdh57/ node02:$PWD
scp -r apache-kylin-2.6.3-bin-cdh57/ node01:$PWD

第二步:三台机器修改kylin配置文件kylin.properties

三台服务器分别修改kylin配置文件kylin.properties

node01服务器修改配置文件

cd /kkb/install/apache-kylin-2.6.3-bin-cdh57/conf/

vim kylin.properties

kylin.metadata.url=kylin_metadata@hbase

kylin.env.hdfs-working-dir=/kylin

kylin.server.mode=query

kylin.server.cluster-servers=node01:7070,node02:7070,node03:7070

kylin.storage.url=hbase

kylin.job.retry=2

kylin.job.max-concurrent-jobs=10

kylin.engine.mr.yarn-check-interval-seconds=10

kylin.engine.mr.reduce-input-mb=500

kylin.engine.mr.max-reducer-number=500

kylin.engine.mr.mapper-input-rows=1000000

 

node02服务器修改配置文件

cd /kkb/install/apache-kylin-2.6.3-bin-cdh57/conf/

vim kylin.properties

kylin.metadata.url=kylin_metadata@hbase

kylin.env.hdfs-working-dir=/kylin

kylin.server.mode=query

kylin.server.cluster-servers=node01:7070,node02:7070,node03:7070

kylin.storage.url=hbase

kylin.job.retry=2

kylin.job.max-concurrent-jobs=10

kylin.engine.mr.yarn-check-interval-seconds=10

kylin.engine.mr.reduce-input-mb=500

kylin.engine.mr.max-reducer-number=500

kylin.engine.mr.mapper-input-rows=1000000

 

node03服务器修改配置文件

cd /kkb/install/apache-kylin-2.6.3-bin-cdh57/conf/

vim kylin.properties

kylin.metadata.url=kylin_metadata@hbase

kylin.env.hdfs-working-dir=/kylin

kylin.server.mode=all

kylin.server.cluster-servers=node01:7070,node02:7070,node03:7070

kylin.storage.url=hbase

kylin.job.retry=2

kylin.job.max-concurrent-jobs=10

kylin.engine.mr.yarn-check-interval-seconds=10

kylin.engine.mr.reduce-input-mb=500

kylin.engine.mr.max-reducer-number=500

kylin.engine.mr.mapper-input-rows=1000000

 

第三步:三台机器配置环境变量

三台机器编辑/etc/profile,添加环境变量

注意:需要将hive的安装文件夹,每一台机器都拷贝

sudo vim  /etc/profile

export JAVA_HOME=/kkb/install/jdk1.8.0_141

export PATH=:$JAVA_HOME/bin:$PATH

export HADOOP_HOME=/kkb/install/hadoop-2.6.0-cdh5.14.2

export PATH=:$HADOOP_HOME/bin:$PATH

export HBASE_HOME=/kkb/install/hbase-1.2.0-cdh5.14.2

export PATH=:$HBASE_HOME/bin:$PATH

export HIVE_HOME=/kkb/install/hive-1.1.0-cdh5.14.2

export PATH=:$HIVE_HOME/bin:$PATH

export HCAT_HOME=/kkb/install/hive-1.1.0-cdh5.14.2

export PATH=:$HCAT_HOME/hcatalog:$PATH

export KYLIN_HOME=/kkb/install/apache-kylin-2.6.3-bin-cdh57

export PATH=:$KYLIN_HOME/bin:$PATH

export dir=/kkb/install/apache-kylin-2.6.3-bin-cdh57/bin

export PATH=$dir:$PATH

export HBASE_CLASSPATH=/kkb/install/hbase-1.2.0-cdh5.14.2

export PATH=:$HBASE_CLASSPATH:$PATH

 

第四步:三台机器启动kylin服务

三台机器执行以下命令启动kylin服务

cd /kkb/soft/apache-kylin-2.6.3-bin-cdh57
bin/kylin.sh start

第五步:node02安装nginx实现请求负载均衡

注意:nginx的安装需要使用root用户来进行安装

在node02服务器上面安装nginx服务,实现请求负载均衡

将nginx的安装包上传到/kkb/soft路径下,然后解压,并对nginx的配置文件进行配置,然后启动nginx服务即可

1、解压nginx压缩吧

cd /kkb/soft/
tar -zxf nginx-1.8.1.tar.gz -C /kkb/install/

2、编译nginx

yum -y install gcc pcre-devel zlib-devel openssl openssl-devel
cd /kkb/install/nginx-1.8.1/
./configure --prefix=/usr/local/nginx
make
make install

3、修改nginx的配置文件

node02执行以下命令修改nginx的配置文件

cd /usr/local/nginx/conf
vim nginx.conf

添加以下内容

在nginx.conf配置文件的最后一个 “}” 上面一行,添加以下内容

upstream kaikeba {
              least_conn;
              server 192.168.52.100:7070 weight=8;
              server 192.168.52.110:7070 weight=7;
              server 192.168.52.120:7070 weight=7;
       }
       server {
              listen 8066;
              server_name localhost;
              location / {
              proxy_pass http://kaikeba;
              }
       }

4、nginx的启动与停止命令

nginx的启动命令,node02执行以下命令启动nginx服务

cd /usr/local/nginx/
sbin/nginx  -c conf/nginx.conf

nginx的停止命令,node02执行以下命令停止nginx服务

cd /usr/local/nginx/
sbin/nginx -s stop

第六步:浏览器界面访问

http://node02:8066/kylin/

访问这个网址,就可以实现负载均衡

4、kylin的入门使用

我们kylin环境安装成功之后,我们就可以在hive当中创建数据库以及数据库表,然后通过kylin来实现数据的查询

第一步:创建hive数据库以及表并加载以下数据

注意,以下两个文件的字段分隔符均为‘\t’。

dept.txt

10 ACCOUNTING 1700
20 RESEARCH 1800
30 SALES 1900
40 OPERATIONS 1700

emp.txt


7369 SMITH CLERK 7902 1980-12-17 800.00 20
7499 ALLEN SALESMAN 7698 1981-2-20 1600.00 300.00 30
7521 WARD SALESMAN 7698 1981-2-22 1250.00 500.00 30
7566 JONES MANAGER 7839 1981-4-2 2975.00 20
7654 MARTIN SALESMAN 7698 1981-9-28 1250.00 1400.00 30
7698 BLAKE MANAGER 7839 1981-5-1 2850.00 30
7782 CLARK MANAGER 7839 1981-6-9 2450.00 10
7788 SCOTT ANALYST 7566 1987-4-19 3000.00 20
7839 KING PRESIDENT 1981-11-17 5000.00 10
7844 TURNER SALESMAN 7698 1981-9-8 1500.00 0.00 30
7876 ADAMS CLERK 7788 1987-5-23 1100.00 20
7900 JAMES CLERK 7698 1981-12-3 950.00 30
7902 FORD ANALYST 7566 1981-12-3 3000.00 20
7934 MILLER CLERK 7782 1982-1-23 1300.00 10

将以上两份文件上传到node03服务器的/kkb/install路径下,然后执行以下命令,创建hive数据库以及数据库表,并加载数据

cd /kkb/install/hive-1.1.0-cdh5.14.2/
bin/beeline

创建数据库并使用该数据库

create database kylin_hive;
use kylin_hive;

(1)创建部门表

create external table if not exists kylin_hive.dept(
deptno int,
dname string,
loc int )
row format delimited fields terminated by '\t';

(2)创建员工表

create external table if not exists kylin_hive.emp(
empno int,
ename string,
job string,
mgr int,
hiredate string,
sal double,
comm double,
deptno int)
row format delimited fields terminated by '\t';

(3)查看创建的表

jdbc:hive2://node03:10000> show tables;
OK

tab_name
dept
emp

(4)向外部表中导入数据导入数据

load data local inpath '/kkb/install/dept.txt' into table kylin_hive.dept;
load data local inpath '/kkb/install/emp.txt' into table kylin_hive.emp;

查询结果

jdbc:hive2://node03:10000> select * from emp;
jdbc:hive2://node03:10000> select * from dept;

第二步:访问kylin浏览器界面,并创建project

直接在浏览器界面访问

http://node02:8066/kylin/login  并登录kylin,用户名  ADMIN,密码KYLIN

点击页面 + 号,来创建工程

file

输入工程名称以及工程描述

file

为工程添加数据源

file

添加数据源表

第三步:为kylin添加models

1、回到models页面

2、添加new models

3、填写model name之后,继续下一步

file

4、选择事实表

这里就选择emp作为事实表

file

5、添加维度表

添加我们的DEPT作为维度表,并选择我们的join方式,以及join连接字段

file

6、选择聚合维度信息

file

7、选择度量信息

file

8、添加分区信息及过滤条件之后“Save”

file

第四步:通过kylin来构建cube

前面我们已经创建了project和我们的models,接下来我们就来构建我们的cube

1、页面添加,创建一个new  cube

2、选择我们的model以及cube name

3、添加我们的自定义维度

file

4、添加统计维度

 

file

5、设置多个分区cube合并信息

因为我们这里是全量统计,不涉及多个分区cube进行合并,所以不用设置历史多个cube进行合并

file

6、高级设置

高级设置我们这里暂时也不做任何设置,后续再单独详细讲解

file file file

7、额外的其他的配置属性,这里也暂时不做配置

file

8、完成,保存配置

file

第五步:构建我们的cube

将我们的cube进行构建

file

file

注意:此步骤非常慢,需耐心等待,如果发生错误,可以展开右边箭头查看日志。

第六步:对我们的数进行查询

前面构建好了我们的cube之后,接下来我们就可以对我们的数据进行分析

SELECT  DEPT.DNAME ,SUM(EMP.SAL) FROM EMP  INNER JOIN DEPT  ON DEPT.DEPTNO = EMP.DEPTNO  GROUP BY DEPT.DNAME

file

我们会发现,数据的查询速度非常快,马上就可以产出结果了,通过kylin的与计算,已经将我们各种可能性的结果都获取到了,我们这里直接就可以得到我们计算完成的结果,所以结果非常快就能计算出来

5、kylin的构建流程

双击下面图片,即可播放PPT浏览查看

file file file file file file

6、cube构建算法

6.1、逐层构建算法

file

我们知道,一个N维的Cube,是由1个N维子立方体、N个(N-1)维子立方体、N*(N-1)/2个(N-2)维子立方体、......、N个1维子立方体和1个0维子立方体构成,总共有2^N个子立方体组成,在逐层算法中,按维度数逐层减少来计算,每个层级的计算(除了第一层,它是从原始数据聚合而来),是基于它上一层级的结果来计算的。比如,[Group by A, B]的结果,可以基于[Group by A, B, C]的结果,通过去掉C后聚合得来的;这样可以减少重复计算;当 0维度Cuboid计算出来的时候,整个Cube的计算也就完成了。

每一轮的计算都是一个MapReduce任务,且串行执行;一个N维的Cube,至少需要N次MapReduce Job。

file

算法优点:

1)此算法充分利用了MapReduce的优点,处理了中间复杂的排序和shuffle工作,故而算法代码清晰简单,易于维护;

2)受益于Hadoop的日趋成熟,此算法非常稳定,即便是集群资源紧张时,也能保证最终能够完成。

算法缺点:

1)当Cube有比较多维度的时候,所需要的MapReduce任务也相应增加;由于Hadoop的任务调度需要耗费额外资源,特别是集群较庞大的时候,反复递交任务造成的额外开销会相当可观;

2)由于Mapper逻辑中并未进行聚合操作,所以每轮MR的shuffle工作量都很大,导致效率低下。

3)对HDFS的读写操作较多:由于每一层计算的输出会用做下一层计算的输入,这些Key-Value需要写到HDFS上;当所有计算都完成后,Kylin还需要额外的一轮任务将这些文件转成HBase的HFile格式,以导入到HBase中去;

总体而言,该算法的效率较低,尤其是当Cube维度数较大的时候。

6.2、快速构建算法

file

也被称作“逐段”(By Segment) 或“逐块”(By Split) 算法,从1.5.x开始引入该算法,该算法的主要思想是,每个Mapper将其所分配到的数据块,计算成一个完整的小Cube 段(包含所有Cuboid)。每个Mapper将计算完的Cube段输出给Reducer做合并,生成大Cube,也就是最终结果。如图所示解释了此流程。

file

与旧算法相比,快速算法主要有两点不同:

1) Mapper会利用内存做预聚合,算出所有组合;Mapper输出的每个Key都是不同的,这样会减少输出到Hadoop MapReduce的数据量,Combiner也不再需要;

2)一轮MapReduce便会完成所有层次的计算,减少Hadoop任务的调配。

7、cube构建的优化

从之前章节的介绍可以知道,在没有采取任何优化措施的情况下,Kylin会对每一种维度的组合进行预计算,每种维度的组合的预计算结果被称为Cuboid。假设有4个维度,我们最终会有24 =16个Cuboid需要计算。

但在现实情况中,用户的维度数量一般远远大于4个。假设用户有10 个维度,那么没有经过任何优化的Cube就会存在210 =1024个Cuboid;而如果用户有20个维度,那么Cube中总共会存在220 =1048576个Cuboid。虽然每个Cuboid的大小存在很大的差异,但是单单想到Cuboid的数量就足以让人想象到这样的Cube对构建引擎、存储引擎来说压力有多么巨大。因此,在构建维度数量较多的Cube时,尤其要注意Cube的剪枝优化(即减少Cuboid的生成)。

7.1、使用衍生维度(derived dimension)

衍生维度用于在有效维度内将维度表上的非主键维度排除掉,并使用维度表的主键(其实是事实表上相应的外键)来替代它们。Kylin会在底层记录维度表主键与维度表其他维度之间的映射关系,以便在查询时能够动态地将维度表的主键“翻译”成这些非主键维度,并进行实时聚合。

file

虽然衍生维度具有非常大的吸引力,但这也并不是说所有维度表上的维度都得变成衍生维度,如果从维度表主键到某个维度表维度所需要的聚合工作量非常大,则不建议使用衍生维度。

7.2、 使用聚合组(Aggregation group)

聚合组(Aggregation Group)是一种强大的剪枝工具。聚合组假设一个Cube的所有维度均可以根据业务需求划分成若干组(当然也可以是一个组),由于同一个组内的维度更可能同时被同一个查询用到,因此会表现出更加紧密的内在关联。每个分组的维度集合均是Cube所有维度的一个子集,不同的分组各自拥有一套维度集合,它们可能与其他分组有相同的维度,也可能没有相同的维度。每个分组各自独立地根据自身的规则贡献出一批需要被物化的Cuboid,所有分组贡献的Cuboid的并集就成为了当前Cube中所有需要物化的Cuboid的集合。不同的分组有可能会贡献出相同的Cuboid,构建引擎会察觉到这点,并且保证每一个Cuboid无论在多少个分组中出现,它都只会被物化一次。

对于每个分组内部的维度,用户可以使用如下三种可选的方式定义,它们之间的关系,具体如下。

1)强制维度(Mandatory),如果一个维度被定义为强制维度,那么这个分组产生的所有Cuboid中每一个Cuboid都会包含该维度。每个分组中都可以有0个、1个或多个强制维度。如果根据这个分组的业务逻辑,则相关的查询一定会在过滤条件或分组条件中,因此可以在该分组中把该维度设置为强制维度。

2)层级维度(Hierarchy),每个层级包含两个或更多个维度。假设一个层级中包含D1,D2…Dn这n个维度,那么在该分组产生的任何Cuboid中, 这n个维度只会以(),(D1),(D1,D2)…(D1,D2…Dn)这n+1种形式中的一种出现。每个分组中可以有0个、1个或多个层级,不同的层级之间不应当有共享的维度。如果根据这个分组的业务逻辑,则多个维度直接存在层级关系,因此可以在该分组中把这些维度设置为层级维度。

file

3)联合维度(Joint),每个联合中包含两个或更多个维度,如果某些列形成一个联合,那么在该分组产生的任何Cuboid中,这些联合维度要么一起出现,要么都不出现。每个分组中可以有0个或多个联合,但是不同的联合之间不应当有共享的维度(否则它们可以合并成一个联合)。如果根据这个分组的业务逻辑,多个维度在查询中总是同时出现,则可以在该分组中把这些维度设置为联合维度。

file

这些操作可以在Cube Designer的Advanced Setting中的Aggregation Groups区域完成,如下图所示。

file

聚合组的设计非常灵活,甚至可以用来描述一些极端的设计。假设我们的业务需求非常单一,只需要某些特定的Cuboid,那么可以创建多个聚合组,每个聚合组代表一个Cuboid。具体的方法是在聚合组中先包含某个Cuboid所需的所有维度,然后把这些维度都设置为强制维度。这样当前的聚合组就只能产生我们想要的那一个Cuboid了。

再比如,有的时候我们的Cube中有一些基数非常大的维度,如果不做特殊处理,它就会和其他的维度进行各种组合,从而产生一大堆包含它的Cuboid。包含高基数维度的Cuboid在行数和体积上往往非常庞大,这会导致整个Cube的膨胀率变大。如果根据业务需求知道这个高基数的维度只会与若干个维度(而不是所有维度)同时被查询到,那么就可以通过聚合组对这个高基数维度做一定的“隔离”。我们把这个高基数的维度放入一个单独的聚合组,再把所有可能会与这个高基数维度一起被查询到的其他维度也放进来。这样,这个高基数的维度就被“隔离”在一个聚合组中了,所有不会与它一起被查询到的维度都没有和它一起出现在任何一个分组中,因此也就不会有多余的Cuboid产生。这点也大大减少了包含该高基数维度的Cuboid的数量,可以有效地控制Cube的膨胀率。

7.3、 并发粒度优化

当Segment中某一个Cuboid的大小超出一定的阈值时,系统会将该Cuboid的数据分片到多个分区中,以实现Cuboid数据读取的并行化,从而优化Cube的查询速度。具体的实现方式如下:构建引擎根据Segment估计的大小,以及参数“kylin.hbase.region.cut”的设置决定Segment在存储引擎中总共需要几个分区来存储,如果存储引擎是HBase,那么分区的数量就对应于HBase中的Region数量。kylin.hbase.region.cut的默认值是5.0,单位是GB,也就是说对于一个大小估计是50GB的Segment,构建引擎会给它分配10个分区。用户还可以通过设置kylin.hbase.region.count.min(默认为1)和kylin.hbase.region.count.max(默认为500)两个配置来决定每个Segment最少或最多被划分成多少个分区。

file

由于每个Cube的并发粒度控制不尽相同,因此建议在Cube Designer 的Configuration Overwrites(上图所示)中为每个Cube量身定制控制并发粒度的参数。假设将把当前Cube的kylin.hbase.region.count.min设置为2,kylin.hbase.region.count.max设置为100。这样无论Segment的大小如何变化,它的分区数量最小都不会低于2,最大都不会超过100。相应地,这个Segment背后的存储引擎(HBase)为了存储这个Segment,也不会使用小于两个或超过100个的分区。我们还调整了默认的kylin.hbase.region.cut,这样50GB的Segment基本上会被分配到50个分区,相比默认设置,我们的Cuboid可能最多会获得5倍的并发量。

7.4、 Row Key优化

Kylin会把所有的维度按照顺序组合成一个完整的Rowkey,并且按照这个Rowkey升序排列Cuboid中所有的行。

设计良好的Rowkey将更有效地完成数据的查询过滤和定位,减少IO次数,提高查询速度,维度在rowkey中的次序,对查询性能有显著的影响。

Row key的设计原则如下:

1)被用作where过滤的维度放在前边。

file

2)基数大的维度放在基数小的维度前边。

file

7.5、增量cube构建

我们前面可以构建全量cube,也可以实现增量cube的构建,就是通过分区表的分区时间字段来进行怎量构建

  1. 更改model

file

2、更改cube

file

file

8、备份以及恢复kylin的元数据信息

Kylin组织它所有的元数据(包括cube descriptions and instances, projects, inverted index description and instances,jobs, tables and dictionaries)作为一个层次的文件系统。

然而,Kylin使用HBase来进行存储,而不是普通的文件系统。

我们可以从Kylin的配置文件kylin.properties中查看到:

## The metadata store in hbase
kylin.metadata.url=kylin_metadata@hbase

表示Kylin的元数据被保存在HBase的kylin_metadata表中。

Kylin自身提供了元数据的备份程序,我们可以执行程序看一下帮助信息:

bin/metastore.sh

usage: metastore.sh backup
metastore.sh fetch DATA
metastore.sh reset
metastore.sh refresh-cube-signature
metastore.sh restore PATH_TO_LOCAL_META
metastore.sh list RESOURCE_PATH
metastore.sh cat RESOURCE_PATH
metastore.sh remove RESOURCE_PATH
metastore.sh clean [--delete true]

备份元数据

bin/metastore.sh backup

恢复元数据

bin/metastore.sh reset

接着,上传备份的元数据到Kylin的元数据中

bin/metastore.sh restore $KYLIN_HOME/meta_backups/meta_xxxx_xx_xx_xx_xx_xx

等待操作成功,用户在页面点击Reload Metadata按钮对元数据缓存进行刷新,即可看到最新的元数据

9、kylin的垃圾清理

当kylin运行一段时间后,有很多数据因为不在使用就变成了垃圾数据,这些数据占据着HDFS HBase等资源,当积累到一定程度会对集群性能产生影响。

清理元数据

清理元数据指从kylin元数据中清理掉无用的资源。比如字典表的快照变得无用了。

步骤:

检查哪些资源可以清理,这一步不会删除任何东西:

bin/metastore.sh clean

这会列出所有可以被清理的资源供用户核对,并不会实际上进行删除。

在上述命令中 添加 --delete true .这样就会清理掉晚一点资源,注意操作前最好备份一下元数据

bin/metastore.sh clean --delete true

清理存储器数据

1. 检查哪些资源需要被清理,这个操作不会删除任何内容:

${KYLIN_HOME}/bin/kylin.sh org.apache.kylin.storage.hbase.util.StorageCleanupJob --delete

false

2. 根据上面的输出结果,挑选一两个资源看看是否是不再需要的。接着,在上面的命令基础上添加“–

delete true”选项,开始执行清理操作,命令执行完成后,中间的HDFS文件和HTables表就被删除了。

${KYLIN_HOME}/bin/kylin.sh org.apache.kylin.storage.hbase.util.StorageCleanupJob --delete

true

10、BI工具集成

http://kylin.apache.org/cn/docs/howto/howto_use_restapi.html

官方文档使用说明

可以与Kylin结合使用的可视化工具很多,例如:

ODBC:与Tableau、Excel、PowerBI等工具集成

JDBC:与Saiku、BIRT等Java工具集成

RestAPI:与JavaScript、Web网页集成

Kylin开发团队还贡献了Zepplin的插件,也可以使用Zepplin来访问Kylin服务。

10.1、JDBC

1)新建项目并导入依赖

<dependencies>
    <dependency>
        <groupId>org.apache.kylin</groupId>
        <artifactId>kylin-jdbc</artifactId>
        <version>2.5.1</version>
    </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>

2)编码

package com.kkb.kylin;

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;

public class KylinJdbc {
    public static void main(String[] args) throws Exception {

        //Kylin_JDBC 驱动
        String KYLIN_DRIVER = "org.apache.kylin.jdbc.Driver";

        //Kylin_URL
        String KYLIN_URL = "jdbc:kylin://node02:8066/kylin_hive";

        //Kylin的用户名
        String KYLIN_USER = "ADMIN";

        //Kylin的密码
        String KYLIN_PASSWD = "KYLIN";

        //添加驱动信息
        Class.forName(KYLIN_DRIVER);

        //获取连接
        Connection connection = DriverManager.getConnection(KYLIN_URL, KYLIN_USER, KYLIN_PASSWD);

        //预编译SQL
        PreparedStatement ps = connection.prepareStatement("SELECT sum(sal) FROM emp group by deptno");

        //执行查询
        ResultSet resultSet = ps.executeQuery();

        //遍历打印
        while (resultSet.next()) {
            System.out.println(resultSet.getInt(1));
        }
    }
}

3)结果展示

11875
3750
9400

11、使用kylin来分析我们Hbase当中的数据

前面我们已经通过flink将数据介入到了hbase当中去了,那么我们接下来就可以通过hive整合hbase,将hbase当中的数据映射到hive表当中来,然后通过kylin来对hive当中的数据进行预分析,实现实时数仓的统计功能

第一步:拷贝hbase的五个jar包到hive的lib目录下

将我们HBase的五个jar包拷贝到hive的lib目录下

hbase的jar包都在/kkb/install/hbase-1.2.0-cdh5.14.2/lib

我们需要拷贝五个jar包名字如下

hbase-client-1.2.0-cdh5.14.2.jar              
hbase-hadoop2-compat-1.2.0-cdh5.14.2.jar
hbase-hadoop-compat-1.2.0-cdh5.14.2.jar 
hbase-it-1.2.0-cdh5.14.2.jar   
hbase-server-1.2.0-cdh5.14.2.jar

我们直接在node03执行以下命令,通过创建软连接的方式来进行jar包的依赖

ln -s /kkb/install/hbase-1.2.0-cdh5.14.2/lib/hbase-client-1.2.0-cdh5.14.2.jar              /kkb/install/hive-1.1.0-cdh5.14.2/lib/hbase-client-1.2.0-cdh5.14.2.jar            

ln -s /kkb/install/hbase-1.2.0-cdh5.14.2/lib/hbase-hadoop2-compat-1.2.0-cdh5.14.2.jar      /kkb/install/hive-1.1.0-cdh5.14.2/lib/hbase-hadoop2-compat-1.2.0-cdh5.14.2.jar            

ln -s /kkb/install/hbase-1.2.0-cdh5.14.2/lib/hbase-hadoop-compat-1.2.0-cdh5.14.2.jar       /kkb/install/hive-1.1.0-cdh5.14.2/lib/hbase-hadoop-compat-1.2.0-cdh5.14.2.jar           

ln -s /kkb/install/hbase-1.2.0-cdh5.14.2/lib/hbase-it-1.2.0-cdh5.14.2.jar     /kkb/install/hive-1.1.0-cdh5.14.2/lib/hbase-it-1.2.0-cdh5.14.2.jar              

ln -s /kkb/install/hbase-1.2.0-cdh5.14.2/lib/hbase-server-1.2.0-cdh5.14.2.jar          /kkb/install/hive-1.1.0-cdh5.14.2/lib/hbase-server-1.2.0-cdh5.14.2.jar   

 第二步:修改hive的配置文件

编辑node03服务器上面的hive的配置文件hive-site.xml添加以下两行配置

cd /kkb/install/hive-1.1.0-cdh5.14.2/conf
vim hive-site.xml

<property>
    <name>hive.zookeeper.quorum</name>
    <value>node01,node02,node03</value>
</property>
<property>
    <name>hbase.zookeeper.quorum</name>
    <value>node01,node02,node03</value>
</property>

第三步:修改hive-env.sh配置文件添加以下配置

cd /kkb/install/hive-1.1.0-cdh5.14.2/conf
vim hive-env.sh

export HADOOP_HOME=/kkb/install/hadoop-2.6.0-cdh5.14.2
export HBASE_HOME=/kkb/install/hbase-1.2.0-cdh5.14.2/
export HIVE_CONF_DIR=/kkb/install/hive-1.1.0-cdh5.14.2/conf

第四步:创建hive表,映射hbase当中的数据

进入hive客户端,创建hive映射表,映射hbase当中的两张表数据

create database hive_hbase;

use hive_hbase;

CREATE external TABLE hive_hbase.data_goods(goodsId int ,goodsName     string ,sellingPrice  string ,productPic    string ,productBrand  string ,productfbl    string ,productNum    string ,productUrl    string ,productFrom   string ,goodsStock    int  ,  appraiseNum   int)

STORED BY 'org.apache.hadoop.hive.hbase.HBaseStorageHandler' WITH SERDEPROPERTIES

("hbase.columns.mapping" = ":key,f1:goodsName    ,f1:sellingPrice ,f1:productPic   ,f1:productBrand ,f1:productfbl   ,f1:productNum   ,f1:productUrl   ,f1:productFrom  ,f1:goodsStock   ,    f1:appraiseNum")

TBLPROPERTIES("hbase.table.name" ="flink:data_goods");

CREATE external TABLE hive_hbase.data_orders(orderId int,orderNo string ,userId int,goodId int ,goodsMoney decimal(11,2) ,realTotalMoney  decimal(11,2) ,payFrom         int           ,province        string        ,createTime      timestamp )

STORED BY 'org.apache.hadoop.hive.hbase.HBaseStorageHandler' WITH SERDEPROPERTIES

("hbase.columns.mapping" = ":key,  f1:orderNo ,  f1:userId        ,   f1:goodId        ,   f1:goodsMoney    ,f1:realTotalMoney,f1:payFrom       ,f1:province,f1:createTime")

TBLPROPERTIES("hbase.table.name" ="flink:data_orders");

第五步:在kylin当中对我们hive的数据进行多维度分析

直接登录kylin的管理界面,对我们hive当中的数据进行多维度分析.

 

Views: 86

Hadoop集群可视化管理- Hue

一、课前准备

  1. 准备好大数据集群,启动所有的服务,例如hadoop,hbase,impala,hiveserver2,mysql等各种服务

二、课堂主题

本堂课主要介绍hue这个图形化的界面工具,以及与其他工具之间的整合使用

三、课堂目标

  1. 实现hue与其他框架的整合使用

四、知识要点

1、hue的基本介绍

HUE=Hadoop User Experience

Hue是一个开源的Apache Hadoop UI系统,由Cloudera Desktop演化而来,最后Cloudera公司将其贡献给Apache基金会的Hadoop社区,它是基于Python Web框架Django实现的。

通过使用Hue我们可以在浏览器端的Web控制台上与Hadoop集群进行交互来分析处理数据,例如操作HDFS上的数据,运行MapReduce Job,执行Hive的SQL语句,浏览HBase数据库等等。

HUE链接

· Site: http://gethue.com/

· Github: https://github.com/cloudera/hue

· Reviews: https://review.cloudera.org

Hue的架构

1571452158221

核心功能

· SQL编辑器,支持Hive, Impala, MySQL, Oracle, PostgreSQL, SparkSQL, Solr SQL, Phoenix…

· 搜索引擎Solr的各种图表

· Spark和Hadoop的友好界面支持

· 支持调度系统Apache Oozie,可进行workflow的编辑、查看

HUE提供的这些功能相比Hadoop生态各组件提供的界面更加友好,但是一些需要debug的场景可能还是需要使用原生系统才能更加深入的找到错误的原因。

HUE中查看Oozie workflow时,也可以很方便的看到整个workflow的DAG图,不过在最新版本中已经将DAG图去掉了,只能看到workflow中的action列表和他们之间的跳转关系,想要看DAG图的仍然可以使用oozie原生的界面系统查看。

1,访问HDFS和文件浏览

2,通过web调试和开发hive以及数据结果展示

3,查询solr和结果展示,报表生成

4,通过web调试和开发impala交互式SQL Query

5,spark调试和开发

7,oozie任务的开发,监控,和工作流协调调度

8,Hbase数据查询和修改,数据展示

9,Hive的元数据(metastore)查询

10,MapReduce任务进度查看,日志追踪

11,创建和提交MapReduce,Streaming,Java job任务

12,Sqoop2的开发和调试

13,Zookeeper的浏览和编辑

14,数据库(MySQL,PostGres,SQlite,Oracle)的查询和展示

一句话总结:Hue是一个友好的界面集成框架,可以集成我们各种学习过的以及将要学习的框架,一个界面就可以做到查看以及执行所有的框架

2、Hue的安装

Hue的安装支持多种方式,包括rpm包的方式进行安装,tar.gz包的方式进行安装以及cloudera manager的方式来进行安装等,我们这里使用tar.gz包的方式来进行安装

第一步:下载Hue的压缩包并上传到linux解压

Hue的压缩包的下载地址:

http://archive.cloudera.com/cdh5/cdh/5/

我们这里使用的是CDH5.14.2这个对应的版本,具体下载地址为

http://archive.cloudera.com/cdh5/cdh/5/hue-3.9.0-cdh5.14.2.tar.gz

下载然后上传到node03服务器的/kkb/soft路径下

cd /kkb/soft
tar -zxvf hue-3.9.0-cdh5.14.2.tar.gz -C  /kkb/install

第二步:编译安装启动

2.1、linux系统安装依赖包:

联网安装各种必须的依赖包

sudo yum install ant asciidoc cyrus-sasl-devel cyrus-sasl-gssapi cyrus-sasl-plain gcc gcc-c++ krb5-devel libffi-devel libxml2-devel libxslt-devel make  mysql mysql-devel openldap-devel python-devel sqlite-devel gmp-devel libffi  gcc gcc-c++ kernel-devel openssl-devel gmp-devel openldap-devel
2.2、开始配置Hue
cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf

vim  hue.ini
#通用配置

[desktop]
    secret_key=jFE93j;2[290-eiw.KEiwN2s3['d;/.q[eIW^y#e=+Iei*@Mn<qW5o
    http_host=node03.kaikeba.com
    is_hue_4=true
    time_zone=Asia/Shanghai
    server_user=hadoop
    server_group=hadoop
    default_user=hadoop
    default_hdfs_superuser=hadoop

#配置使用mysql作为hue的存储数据库,大概在hue.ini的587行左右

[[database]]
    engine=mysql
    host=node03.kaikeba.com
    port=3306
    user=root
    password=123456
    name=hue
2.3、创建mysql数据库

创建hue数据库

create database hue default character set utf8 default collate utf8_general_ci;

注意:实际工作中,还需要为hue这个数据库创建对应的用户,并分配权限,我这就不创建了,所以下面这一步不用执行了

grant all on hue.* to 'hue'@'%' identified by 'hue';
2.4、准备进行编译

node03服务器执行以下命令准备进行编译

cd /kkb/install/hue-3.9.0-cdh5.14.2
make apps

注意:如果编译失败,请参照第一步重新解压

2.5、linux系统添加普通用户hue

为了方便也可以使用已经存在的hadoop用户来操作hue。

而工作中往往会单独创建一个用户hue:需要的话可在node03执行以下命令,创建普通用户hue

sudo useradd hue
sudo passwd hue
2.6、启动hue进程

node03执行以下命令启动hue

cd /kkb/install/hue-3.9.0-cdh5.14.2
sudo build/env/bin/supervisor
2.7、页面访问

http://node03:8888

第一次访问的时候,需要设置管理员用户和密码。我们这里的管理员的用户名与密码尽量保持与我们安装hadoop的用户名和密码一致,我们安装hadoop的用户名与密码分别是hadoop 123456,初次登录使用hadoop,密码为123456。

访问页面异常

如果忘记在hue.ini文件中修改元数据库引擎由sqlite改为mysql,访问页面将会遇到如下错误

File "/usr/local/lib/python2.7/site-packages/django/db/backends/sqlite3/base.py", line 323, in execute
  return Database.Cursor.execute(self, query, params)
OperationalError: attempt to write a readonly database

解决办法,检查hue.ini配置是否将元数据库引擎由sqlite改为mysql,以及其他配置是否正确:

[[database]]
engine=mysql        #默认是sqlite3作为元数据库,这里改为mysql
host=192.168.186.36 #<mysql所在服务器>
port=3306           #<mysql端口,一般就是3306>
user=hive           #<mysql用户名>
password=hive1234   #<mysql用户密码>
name=hue            #<数据库名称,新数据库,专门用于hue,里面现在没有任何表>

完成以上的这个配置,启动Hue,通过浏览器访问,仍然会发生错误

DatabaseError: (1146, "Table 'hue.desktop_settings' doesn't exist")

原因是mysql数据没有被初始化,解决办法, 初始化数据库(执行一次即可):

sudo build/env/bin/hue syncdb
sudo build/env/bin/hue migrate

同步数据库的时候会询问是否创建用户,可以选择创建hue用户,并设计密码为hue。

这一步是可选的,因为就算不创建,登陆的时候也需要创建,一般使用Linux中一样的用户名和密码即可。(我用的是hadoop用户)

然后再启动supervisor服务(前台运行,如果需要后台运行可以添加-d选项 )

sudo build/env/bin/supervisor

3、hue与其他框架的集成

如果登陆后仍然由报错,很可能是因为hue和hadoop的hdfs和yarn的集成没有设置好。

3.1、hue与hadoop的HDFS以及yarn集成

第一步:更改所有hadoop节点的core-site.xml配置

记得更改完core-site.xml之后一定要重启hdfs与yarn集群

三台机器更改core-site.xml

<property>
    <name>hadoop.proxyuser.hadoop.hosts</name>
    <value>*</value>
</property>

<property>
    <name>hadoop.proxyuser.hadoop.groups</name>
    <value>*</value>
</property> 
第二步:更改所有hadoop节点的hdfs-site.xml

所有服务器更改hdfs-site.xml添加以下配置

<property>
    <name>dfs.webhdfs.enabled</name>
    <value>true</value>
</property>
第三步:重启hadoop集群

在node01机器上面执行以下命令

cd /kkb/install/hadoop-2.6.0-cdh5.14.2

sbin/stop-dfs.sh
sbin/start-dfs.sh
sbin/stop-yarn.sh
sbin/start-yarn.sh
第四步:停止hue的服务,并继续配置hue.ini

停止hue的服务,然后进入到以下路径,重新配置hue.ini这个配置文件(880行左右)

cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf

vim hue.ini

#配置我们的hue与hdfs集成]]
[[hdfs_clusters]]
    [[[default]]]
        fs_defaultfs=hdfs://node01.kaikeba.com:8020
        webhdfs_url=http://node01.kaikeba.com:50070/webhdfs/v1
        hadoop_hdfs_home=/kkb/install/hadoop-2.6.0-cdh5.14.2
        hadoop_bin=/kkb/install/hadoop-2.6.0-cdh5.14.2/bin
        hadoop_conf_dir=/kkb/install/hadoop-2.6.0-cdh5.14.2/etc/hadoop

#配置我们的hue与yarn集成

[[yarn_clusters]]

    [[[default]]]

      resourcemanager_host=node01
      resourcemanager_port=8032
      submit_to=True
      resourcemanager_api_url=http://node01:8088
      history_server_api_url=http://node01:19888

注意:如果是hadoop 3的版本,那么webhdfs_url端口一般需要修改成9870

webhdfs_url=http://node01.kaikeba.com:9870/webhdfs/v1

配置完成之后重新启动hue的服务

node03执行以下命令进行重新启动hue的服务

$ cd /kkb/install/hue-3.9.0-cdh5.14.2/
$ sudo build/env/bin/supervisor

如果出现8888端口占用,需要先kill掉再启动服务

3.2、配置hue与hive集成

如果需要配置hue与hive的集成,我们需要启动hive的metastore服务以及hiveserver2服务(impala需要hive的metastore服务,hue需要hvie的hiveserver2服务)

更改hue的配置hue.ini

停止hue的服务,然后重新编辑修改hue.ini这个配置文件

修改hue.ini

cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf
vim hue.ini

[beeswax]
    hive_server_host=node03.kaikeba.com
    hive_server_port=10000
    hive_conf_dir=/kkb/install/hive-1.1.0-cdh5.14.2/conf
    server_conn_timeout=120
    auth_username=hadoop
    auth_password=123456

[metastore]

  #允许使用hive创建数据库表等操作
  enable_new_create_table=true
启动hive的metastore服务

去node03机器上启动hive的metastore以及hiveserver2服务

cd /kkb/install/hive-1.1.0-cdh5.14.2/conf
nohup bin/hive --service metastore &
nohup bin/hive --service hiveserver2 &

重新启动hue,然后就可以通过浏览器页面操作hive了

node03执行以下命令进行重新启动hue的服务

cd /kkb/install/hue-3.9.0-cdh5.14.2/
sudo build/env/bin/supervisor
3.3、配置hue与impala的集成

停止hue的服务进程

修改hue.ini配置文件

cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf
vim hue.ini

[impala]

  server_host=node03
  server_port=21050
  impala_conf_dir=/etc/impala/conf

然后node03执行以下命令,重新启动hue的服务即可

cd /kkb/install/hue-3.9.0-cdh5.14.2/
sudo build/env/bin/supervisor

3.4、配置hue与mysql的集成

找到databases 这个选项,将这个选项下面的mysql注释给打开,然后配置mysql即可,大概在1547行

停止hue的服务,然后修改hue.ini

cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf
vim hue.ini

[[[mysql]]]
    nice_name="My SQL DB"
    engine=mysql
    host=node03.kaikeba.com
    port=3306
    user=root
    password=123456
    options={"init_command":"set names utf8;SET CHARACTER SET utf8;SET character_set_connection=utf8;"}

更改完了配置,重新启动hue的服务

cd /kkb/install/hue-3.9.0-cdh5.14.2/
build/env/bin/supervisor

3.5、配置hue与hbase的集成

第一步:修改hue.ini

如果hue已经启动,需要先停止hue的服务,然后继续修改hue的配置文件hue.ini

cd /kkb/install/hue-3.9.0-cdh5.14.2/desktop/conf

vim hue.ini

[hbase]
  hbase_clusters=(Cluster|node01:9090)
  hbase_conf_dir=/kkb/install/hbase-1.2.0-cdh5.14.2/conf
第二步:启动hbase的thrift server服务

第一台机器执行以下命令启动hbase的thriftserver

cd /kkb/install/hbase-1.2.0-cdh5.14.2

bin/start-hbase.sh 
bin/hbase-daemon.sh start thrift  

检查hbase主节点状态: http://hadoop101:16010/master-status

第三步:启动hue

第三台机器执行以下命令启动hue

cd /kkb/install/hue-3.9.0-cdh5.14.2/

build/env/bin/supervisor
第四步:页面访问

http://node03:8888/hue/

image-20210613184533712

Views: 99

高速Hive查询引擎 – Impala

一、课前准备

安装好hive以及hadoop运行环境,并正常启动hadoop以及hive的

二、课堂主题

实现impala集群环境正常安装,并掌握impala的基本语法

三、课堂目标

熟练使用impala的语法

四、知识要点

离线任务处理流程概述

离线任务处理流程

由于大部分的软件框架,CDH都提供了压缩包的安装方式,但是由于impala有部分代码使用C++编写,所以impala在安装包的选择上面,cloudera公司没有提供tar包的安装方式,只提供了rpm的安装方式,我们可以通过下载rpm包来进行安装。注意:rpm包是linux操作系统上面的一种安装压缩包

1、 impala的概述

imala基本介绍

impala是cloudera提供的一款高效率的sql查询工具,提供实时的查询效果,官方测试性能比hive快10到100倍,其sql查询比sparkSQL还要更加快速,号称是当前大数据领域最快的查询sql工具,impala是参照谷歌的新三篇论文(Caffeine、Pregel、Dremel)当中的Dremel实现而来,其中旧三篇论文分别是(BigTable,GFS,MapReduce)分别对应我们即将学的HBase和已经学过的HDFS以及MapReduce

impala是基于hive并使用内存进行计算,兼顾数据仓库,具有实时,批处理,多并发等优点

impala与hive的关系

impala是基于hive的大数据分析查询引擎,直接使用hive的元数据库metadata,意味着impala元数据都存储在hive的metastore当中,并且impala兼容hive的绝大多数sql语法。所以需要安装impala的话,必须先安装hive,保证hive安装成功,并且还需要启动hive的metastore服务

impala的优点

1、impala比较快,非常快,特别快,因为所有的计算都可以放入内存当中进行完成,只要你内存足够大

2、摈弃了MR的计算,改用C++来实现,有针对性的硬件优化

3、具有数据仓库的特性,对hive的原有数据做数据分析

4、支持ODBC,jdbc远程访问

impala的缺点:

1、基于内存计算,对内存依赖性较大

2、改用C++编写,意味着维护难度增大

3、基于hive,与hive共存亡,紧耦合

4、稳定性不如hive,不存在数据丢失的情况

impala的架构以及查询计划

img

Impala的架构模块:

  • impala-server

    • 启动的守护进程,执行我们的查询计划 从节点,官方建议与所有的datanode装在一起,可以通过hadoop的短路读取特性实现数据的快速查询
  • impala-statestore

    • 状态存储区 主节点
  • impalas-catalog

    • 元数据管理区 主节点

查询执行

impalad分为frontend和backend两个层次, frondend用java实现(通过JNI嵌入impalad), 负责查询计划生成, 而backend用C++实现, 负责查询执行。

frontend**生成查询计划分为两个阶段:**

(1)生成单机查询计划,单机执行计划与关系数据库执行计划相同,所用查询优化方法也类似。

(2)生成分布式查询计划。 根据单机执行计划, 生成真正可执行的分布式执行计划,降低数据移动, 尽量把数据和计算放在一起。

![http://www.aboutyun.com/data/attachment/forum/201507/27/134940kb7luw6ix8in3w3s.png](https://imgs.delucia.cn/imgs/2021/clip_image004 - 副本.jpg)

上图是SQL查询例子, 该SQL的目标是在三表join的基础上算聚集, 并按照聚集列排序取topN。

impala的查询优化器支持代价模型: 利用表和分区的cardinality,每列的distinct值个数等统计数据, impala可估算执行计划代价, 并生成较优的执行计划。 上图左边是frontend查询优化器生成的单机查询计划, 与传统关系数据库不同, 单机查询计划不能直接执行, 必须转换成如图右半部分所示的分布式查询计划。 该分布式查询计划共分成6个segment(图中彩色无边框圆角矩形), 每个segment是可以被单台服务器独立执行的计划子树。

![img](https://imgs.delucia.cn/imgs/2021/clip_image006 - 副本.jpg)

2、impala的安装环境准备

需要提前安装好hadoop,hive,这两个框架,并且hive需要将hive的安装包,拷贝到所有的服务器上面都保存一份,因为impala需要引用hive的安装目录下面的一些依赖的jar包

3、下载impala的所有依赖包

由于impala没有提供tar包供我们进行安装,只提供了rpm包,所以我们在安装impala的时候,需要使用rpm包来进行安装,rpm包只有cloudera公司提供了,所以我们去cloudera公司网站进行下载rpm包即可,但是另外一个问题,impala的rpm包依赖非常多的其他的rpm包,可以一个个的将依赖找出来,也可以将所有的rpm包下载下来,制作成我们本地yum源来进行安装。我们这里就选择制作我们本地的yum源来进行安装,所以首先我们需要下载到所有的rpm包,下载地址如下

http://archive.cloudera.com/cdh5/repo-as-tarball/5.14.2/cdh5.14.2-centos7.tar.gz

下载好了之后,保留下,留作备用

将我们下载好的压缩包,上传到node03服务器的/kkb/soft路径下,并进行解压

cd /kkb/soft
tar -zxvf cdh5.14.2-centos7.tar.gz

4、制作本地yum源

镜像源是centos当中下载相关软件的地址,我们可以通过制作我们自己的镜像源指定我们去哪里下载impala的rpm包,这里我们使用httpd这个软件来作为服务端,启动httpd的服务来作为我们镜像源的下载地址

这里我们选用第三台机器作为镜像源的服务端

node03机器上执行以下命令

sudo yum  -y install httpd
sudo service httpd start

cd /etc/yum.repos.d
sudo vim localimp.repo 

[localimp]
name=localimp
baseurl=http://node03/cdh5.14.2/
gpgcheck=0
enabled=1

创建apache httpd的读取链接

sudo ln -s /kkb/soft/cdh/5.14.2 /var/www/html/cdh5.14.2

页面访问本地yum源,出现这个界面表示本地yum源制作成功

http://node03/cdh5.14.2

如果能够正常访问到文件浏览页面,证明我们的本地yum源安装成功

将制作好的localimp配置文件发放到所有需要安装impala的节点上去

node03执行以下命令进行分发
cd /etc/yum.repos.d/

sudo scp localimp.repo  node02:$PWD
sudo scp localimp.repo  node01:$PWD

5、开始安装impala

安装规划

服务名称 node01 node02 node03
impala-catalog 不安装 不安装 安装
impala-state-store 不安装 不安装 安装
impala-server 安装 安装 安装
#主节点node03执行以下命令进行安装

sudo yum  install  impala -y
impala-server
impala-state-store
impala-catalog
impala-shell

#从节点node01与node02安装以下服务
sudo yum install impala-server -y

6、所有节点配置impala

第一步:修改hive-site.xml

node03机器修改hive-site.xml内容如下

hive-site.xml配置

vim /kkb/install/hive-1.1.0-cdh5.14.2/conf/hive-site.xml
添加以下三个配置属性

<property>
    <name>hive.server2.thrift.bind.host</name>
    <value>node03.hadoop.com</value>
</property>
 <property>
     <name>hive.metastore.uris</name>
     <value>thrift://node03.kaikeba.com:9083</value>
 </property>
<property>
    <name>hive.metastore.client.socket.timeout</name>
    <value>3600</value>
 </property>

第二步:将hive的安装包发送到node02与node01机器上

在node03机器上面执行

cd /kkb/install/

scp -r hive-1.1.0-cdh5.14.2/ node02:$PWD
scp -r hive-1.1.0-cdh5.14.2/ node01:$PWD

第三步:node03启动hive的metastore服务

启动hive的metastore服务

node03机器启动hive的metastore服务

cd  /kkb/install/hive-1.1.0-cdh5.14.2

nohup bin/hive --service metastore &
nohup bin/hive -- service hiveserver2 &

注意:一定要保证mysql的服务正常启动,否则metastore的服务不能够启动

第四步:所有hadoop节点修改hdfs-site.xml添加以下内容

所有节点创建文件夹

sudo mkdir -p /var/run/hdfs-sockets

修改所有节点的hdfs-site.xml添加以下配置,修改完之后重启hdfs集群生效

vim  /kkb/install/hadoop-2.6.0-cdh5.14.2/etc/hadoop/hdfs-site.xml

<property>
    <name>dfs.client.read.shortcircuit</name>
    <value>true</value>
</property>

<property>
     <name>dfs.domain.socket.path</name>
     <value>/var/run/hdfs-sockets/dn</value>
</property>

<property>
    <name>dfs.client.file-block-storage-locations.timeout.millis</name>
    <value>10000</value>
</property>
<property>
     <name>dfs.datanode.hdfs-blocks-metadata.enabled</name>
     <value>true</value>
</property>

三台机器执行以下命令给文件夹授权

sudo  chown  -R  hadoop:hadoop   /var/run/hdfs-sockets/

第五步:重启hdfs

重启hdfs文件系统

node01服务器上面执行以下命令

cd /kkb/install/hadoop-2.6.0-cdh5.14.2/

sbin/stop-dfs.sh
sbin/start-dfs.sh

第六步:创建hadoop与hive的配置文件的连接

impala的配置目录为 /etc/impala/conf

这个路径下面需要把core-site.xmlhdfs-site.xml以及hive-site.xml拷贝到这里来,但是我们这里使用软连接的方式会更好

所有节点执行以下命令创建链接到impala配置目录下来

sudo ln -s /kkb/install/hadoop-2.6.0-cdh5.14.2/etc/hadoop/core-site.xml /etc/impala/conf/core-site.xml

sudo ln -s /kkb/install/hadoop-2.6.0-cdh5.14.2/etc/hadoop/hdfs-site.xml /etc/impala/conf/hdfs-site.xml

sudo ln -s /kkb/install/hive-1.1.0-cdh5.14.2/conf/hive-site.xml /etc/impala/conf/hive-site.xml

第七步:修改impala的配置文件

所有节点修改impala默认配置

所有节点更改impala默认配置文件以及添加mysql的连接驱动包

sudo vim /etc/default/impala

IMPALA_CATALOG_SERVICE_HOST=node03
IMPALA_STATE_STORE_HOST=node03

所有节点创建mysql的驱动包的软连接

sudo mkdir -p /usr/share/java
sudo ln -s /kkb/install/hive-1.1.0-cdh5.14.2/lib/mysql-connector-java-5.1.38.jar /usr/share/java/mysql-connector-java.jar
所有节点修改bigtop的java路径

修改bigtop的java_home路径

sudo vim /etc/default/bigtop-utils

export JAVA_HOME=/kkb/install/jdk1.8.0_141

第八步:启动impala服务

启动impala服务

主节点node03启动以下三个服务进程

sudo service impala-state-store start
sudo service impala-catalog start
sudo service impala-server start

从节点启动node01与node02启动impala-server

sudo service  impala-server  start

三台机器可以通过以下命令,查看impala进程是否存在

ps -ef | grep impala

注意:启动之后所有关于impala的日志默认都在/var/log/impala 这个路径下,node03机器上面应该有三个进程,node02与node01机器上面只有一个进程,如果进程个数不对,去对应目录下查看报错日志

浏览器页面访问:

访问impalad的管理界面

http://node03:25000/

访问statestored的管理界面

http://node03:25010/

访问catalog的管理界面

http://node03:25020

7、impala的使用

执行impala-shell即可进入交互界面

$ impala-shell

Starting Impala Shell without Kerberos authentication
Connected to node03:21000
Server version: impalad version 2.11.0-cdh5.14.2 RELEASE (build ed85dce709da9557aeb28be89e8044947708876c)
**********************************************************************
Welcome to the Impala shell.
(Impala Shell v2.11.0-cdh5.14.2 (ed85dce) built on Tue Mar 27 13:39:48 PDT 2018)
You can run a single query from the command line using the '-q' option.
**********************************************************************
[node03:21000] >

exit;退出交互界面

1、impala-shell语法

1.1、impala-shell的外部命令参数语法

不需要进入到impala-shell交互命令行当中即可执行的命令参数

impala-shell后面执行的时候可以带很多参数:

-h 查看帮助文档

impala-shell -h

-r 刷新整个元数据,数据量大的时候,比较消耗服务器性能

impala-shell -r

-v 查看对应版本

impala-shell -v -V

-f 执行查询文件

cd /kkb/install

vim impala-shell.sql

select * from course.score;

通过-f 参数来执行执行的查询文件

impala-shell -f impala-shell.sql

-p 显示查询计划

impala-shell -f impala-shell.sql -p

9.1.2、impala-shell的内部命令行参数语法

进入impala-shell命令行之后可以执行的语法

help命令

帮助文档

connect命令

connect hostname 连接到某一台机器上面去执行

refresh 命令

refresh dbname.tablename 增量刷新,刷新某一张表的元数据,主要用于刷新hive当中数据表里面的数据改变的情况

refresh course.score;

invalidate metadata 命令:

invalidate  metadata
全量刷新,性能消耗较大,主要用于hive当中新建数据库或者数据库表的时候来进行刷新

explain 命令:

用于查看sql语句的执行计划

explain select * from course.score;

explain的值可以设置成0,1,2,3等几个值,其中3级别是最高的,可以打印出最全的信息

set explain_level=3;

profile命令:

执行sql语句之后执行,可以打印出更加详细的执行步骤,

主要用于查询结果的查看,集群的调优等

select * from course.score;

profile;

注意:在hive窗口当中插入的数据或者新建的数据库或者数据库表,在impala当中是不可直接查询到的,需要刷新数据库,在impala-shell当中插入的数据,在impala当中是可以直接查询到的,不需要刷新数据库,其中使用的就是catalog这个服务的功能实现的,catalog是impala1.2版本之后增加的模块功能,主要作用就是同步impala之间的元数据

2、创建数据库

impala-shell进入到impala的交互窗口

2.1、查看所有数据库

show databases;

2.2、创建与删除数据库

创建数据库

CREATE DATABASE IF NOT EXISTS mydb1;
drop database  if exists  mydb;

创建数据库表并指定数据库表数据存放hdfs的位置(与hive建表语法类似)

hdfs dfs -mkdir -p /input/impala

create  external table  t3(id int ,name string ,age int )  row  format  delimited fields terminated  by  '\t' location  '/input/impala/external';

3、 创建数据库表

创建student表

CREATE TABLE IF NOT EXISTS mydb1.student (name STRING, age INT, contact INT );

创建employ表

create table employee (Id INT, name STRING, age INT,address STRING, salary BIGINT);
3.1、 数据库表中插入数据
insert into employee (ID,NAME,AGE,ADDRESS,SALARY)VALUES (1, 'Ramesh', 32, 'Ahmedabad', 20000 );
insert into employee values (2, 'Khilan', 25, 'Delhi', 15000 );
Insert into employee values (3, 'kaushik', 23, 'Kota', 30000 );
Insert into employee values (4, 'Chaitali', 25, 'Mumbai', 35000 );
Insert into employee values (5, 'Hardik', 27, 'Bhopal', 40000 );
Insert into employee values (6, 'Komal', 22, 'MP', 32000 );

数据的覆盖

Insert overwrite employee values (1, 'Ram', 26, 'Vishakhapatnam', 37000 );
执行覆盖之后,表中只剩下了这一条数据了

另外一种建表语句

create table customer as select * from employee;
3.2、 数据的查询
select * from employee;

select name,age from employee;
3.3、 删除表
DROP table  mydb1.employee;
3.4、 清空表数据
truncate  employee;
3.5、 查看视图数据
select * from employee_view;

4、 order by语句

基础语法

select * from table_name ORDER BY col_name [ASC|DESC] [NULLS FIRST|NULLS LAST]
Select * from employee ORDER BY id asc;

5、group by 语句

Select name, sum(salary) from employee Group BY name;

6、 having 语句

基础语法

select * from table_name ORDER BY col_name [ASC|DESC] [NULLS FIRST|NULLS LAST]

按年龄对表进行分组,并选择每个组的最大工资,并显示大于20000的工资

select max(salary) from employee group by age having max(salary) > 20000;

7、 limit语句

select * from employee order by id limit 4;

8、impala当中的数据表导入几种方式

第一种方式,通过load hdfs的数据到impala当中去

create table user(id int ,name string,age int ) row format delimited fields terminated by "\t";

准备数据user.txt并上传到hdfs的 /user/impala路径下去

1   hello   15
2   zhangsan    20
3   lisi    30
4   wangwu  50

加载数据

load data inpath '/user/impala/' into table user;

查询加载的数据

select  *  from  user;

如果查询不不到数据,那么需要刷新一遍数据表

refresh  user;

第二种方式:

create  table  user2   as   select * from  user;

第三种方式:

insert into 不推荐使用 因为会产生大量的小文件

千万不要把impala当做一个数据库来使用

第四种:

insert  into  select  用的比较多

9、impala的java开发

在实际工作当中,因为impala的查询比较快,所以可能有会使用到impala来做数据库查询的情况,我们可以通过java代码来进行操作impala的查询

第一步:导入jar包
  <repositories>
        <repository>
            <id>cloudera</id>
            <url>https://repository.cloudera.com/artifactory/cloudera-repos/</url>
        </repository>
        <repository>
            <id>central</id>
            <url>http://repo1.maven.org/maven2/</url>
            <releases>
                <enabled>true</enabled>
            </releases>
            <snapshots>
                <enabled>false</enabled>
            </snapshots>
        </repository>
    </repositories>
    <dependencies>
        <dependency>
            <groupId>org.apache.hadoop</groupId>
            <artifactId>hadoop-common</artifactId>
            <version>2.6.0-cdh5.14.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hive</groupId>
            <artifactId>hive-common</artifactId>
            <version>1.1.0-cdh5.14.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hive</groupId>
            <artifactId>hive-metastore</artifactId>
            <version>1.1.0-cdh5.14.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hive</groupId>
            <artifactId>hive-service</artifactId>
            <version>1.1.0-cdh5.14.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hive</groupId>
            <artifactId>hive-jdbc</artifactId>
            <version>1.1.0-cdh5.14.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.hive</groupId>
            <artifactId>hive-exec</artifactId>
            <version>1.1.0-cdh5.14.2</version>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.thrift/libfb303 -->
        <dependency>
            <groupId>org.apache.thrift</groupId>
            <artifactId>libfb303</artifactId>
            <version>0.9.0</version>
            <type>pom</type>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.apache.thrift/libthrift -->
        <dependency>
            <groupId>org.apache.thrift</groupId>
            <artifactId>libthrift</artifactId>
            <version>0.9.0</version>
            <type>pom</type>
        </dependency>
        <dependency>
            <groupId>org.apache.httpcomponents</groupId>
            <artifactId>httpclient</artifactId>
            <version>4.2.5</version>
        </dependency>
        <dependency>
            <groupId>org.apache.httpcomponents</groupId>
            <artifactId>httpcore</artifactId>
            <version>4.2.5</version>
        </dependency>
    </dependencies>
第二步:impala的java代码查询开发
public class ImpalaJdbc {
     public static void main(String[] args) throws Exception {
     //定义连接驱动类,以及连接url和执行的sql**语句     
     String driver = "org.apache.hive.jdbc.HiveDriver";
     String driverUrl = "jdbc:hive2://192.168.52.120:21050/mydb1;auth=noSasl";
     String sql = "select * from student";

     //通过反射加载数据库连接驱动*    
     Class.forName(driver);
     Connection connection = DriverManager.getConnection(driverUrl);
     PreparedStatement preparedStatement = connection.prepareStatement(sql);
     ResultSet resultSet = preparedStatement.executeQuery();
     //通过查询,得到数据一共有多少列     
     int col = resultSet.getMetaData().getColumnCount();
     //遍历结果集     
     while (resultSet.next()){
         for(int i=1;i<=col;i++){
             System.out.print(resultSet.getString(i)+"\t");
         }
         System.out.print("\n");
     }
     preparedStatement.close();
     connection.close();
 }
 }

Views: 82

Presto分布式SQL查询引擎

一、课前准备

  1. jdk版本要求:Java 8 Update 151 or higher (8u151+), 64-bit
  2. 安装好hadoop集群
  3. 安装好hive

二、课堂主题

  1. 介绍presto
  2. presto架构
  3. prsto安装部署
  4. presto使用

三、课堂目标

  1. 理解presto
  2. 独立完成presto安装部署
  3. 使用presto

四、知识要点

1. Presto是什么?

  • Hadoop提供了大数据存储与计算的一整套解决方案;但是它采用的是MapReduce计算框架,只适合离线和批量计算,无法满足快速实时的Ad-Hoc查询计算的性能要求

  • Hive使用MapReduce作为底层计算框架,是专为批处理设计的。但随着数据越来越多,使用Hive进行一个简单的数据查询可能要花费几分到几小时,显然不能满足交互式查询的需求。

  • Facebook于2012年秋开始开发了Presto,每日查询数据量在1PB级别。Facebook称Presto的性能比Hive要快上10倍多。2013年Facebook正式宣布开源Presto。

  • Presto是apache下开源的==OLAP的分布式SQL查询引擎==,数据量支持从GB到PB级别的数据量的查询,并且查询时,能做到秒级查询。

  • 另外,Presto虽然可以解析SQL,但它并非是标准的数据库;不能替代如MySQL、PostgreSQL、Oracle关系型数据库,不是用于处理OLTP的

  • presto是利用分布式查询,高效的对海量数据进行查询;

  • presto可以用来查询hdfs上的海量数据;但是,presto不仅仅可以用来查询hdfs的数据,它还被设计成能够对很多其他的数据源的数据做查询;

  • 比如数据源有HDFS、Hive、Druid、Kafka、kudu、MySQL、Redis等;下图是Presto 0.237支持的数据源

    image-20200628160739942

2. Presto架构

  • Presto查询引擎是一个Master-Slave的架构,Coordinator是主,worker是从;
  • 一个presto集群,由一个Coordinator节点,一个Discovery Server节点(通常内嵌于Coordinator节点中),多个Worker节点组成
    • Coordinator负责接收查询请求、解析SQL语句、生成执行计划、任务调度给Worker节点执行、worker管理。
    • Worker节点是工作节点;负责实际执行查询任务Task;Worker节点启动后向Discovery Server服务注册;Coordinator从Discovery Server获得可以正常工作的Worker节点。
  • Presto CLI提交查询到Coordinator
  • catalog表示数据源;每个catalog包含Connector及Schema
    • 其中Connector是数据源的适配器;presto通过Connector与不同的数据源(如Redis、Hive、Kafka)连接;如果配置了Hive Connector,需要配置一个Hive MetaStore服务为Presto提供Hive元信息,Worker节点与HDFS交互读取数据。
    • Schema类似于MySQL中的数据库的概念;Schema中又包含Table,类似于MySQL中的表

3. Presto特点

1. 优点

  • 高性能:Presto基于内存计算,减少数据的落盘,计算更快;轻量快速,支持近乎实时的查询
  • 多数据源:通过配置不同的Connector,presto可以连接不同的数据源,所以可以将来自不同数据源的表进行连接查询
  • 支持SQL:完全支持ANSI SQL,并提供了sql shell命令行工具
  • 扩展性:可以根据实际的需要,开发特定的数据源的Connector,从而可以sql查询此数据元的数据

2. 缺点

  • 虽然Presto是基于内存做计算;但是数据量大时,数据并非全部存储在内存中;
    • 比如Presto可针对PB级别的数据做计算,但Presto并非将所有数据全部存储在内存中,不同场景有不同做法;
    • 比如count, avg等聚合运算,会读部分数据,计算,在清理内存;再读数据再计算、清理内存;所以占据内存并不是很高;
    • 但是如果做join操作,中间可能会产生大量的临时数据,造成执行速度变慢;join时,hive的数据反而更快些。所以如果join的话,建议在hive中,先进行join生成宽表,再使用presto查询此宽表数据

3. presto与impala对比

  • impala性能比presto稍好
  • 但是,impala只能对接hive;而presto能对接很多种类的数据源

4. 安装部署Presto

官网地址:https://prestodb.io/

github地址

presto集群规划

主机名 角色
node01 coordinator
node02 worker
node03 worker

1. 安装部署Presto Server

presto要求

image-20200629145103059

确认python版本是2.4+

python -V

image-20200629160952988

确认java版本是8u151+;若如下图,是151之前的版本,安装presto时,需要特殊处理

image-20200629161305214

1. 下载安装包

https://repo1.maven.org/maven2/com/facebook/presto/presto-server/0.237/presto-server-0.237.tar.gz

然后将tar.gz包上传到node01的/kkb/soft目录

2. 解压

cd /kkb/soft/
tar -xzvf presto-server-0.237.tar.gz -C /kkb/install/

3. 配置JAVA

若java版本低于8u151,那么需要上传8u151+的版本压缩包到/kkb/soft;若不低于,则跳过此步骤

解压

cd /kkb/soft/
tar -xzvf jdk-8u251-linux-x64.tar.gz -C /kkb/install/
cd /kkb/install/
scp -r jdk1.8.0_251/ node02:$PWD
scp -r jdk1.8.0_251/ node03:$PWD

presto-server目录创建软链接

ln -s presto-server-0.237/ presto

指定presto使用的java版本(3个节点都要修改)

vim /kkb/install/presto/bin/launcher

添加如下内容

PATH=/kkb/install/jdk1.8.0_251/bin:$PATH
java -version

注意:需要加在exec "$(dirname "$0")/launcher.py" "$@"之前

image-20200629174956052

3. 创建相关目录

创建存储数据文件夹;presto将存储log及其他数据到此目录

cd /kkb/install
cd presto
mkdir data

创建存储配置文件的文件夹

mkdir etc

4. 添加JVM配置文件

etc目录下添加jvm.config配置文件

cd /kkb/install/presto/etc
vim jvm.config

内容如下

-server
-Xmx16G
-XX:+UseG1GC
-XX:G1HeapRegionSize=32M
-XX:+UseGCOverheadLimit
-XX:+ExplicitGCInvokesConcurrent
-XX:+HeapDumpOnOutOfMemoryError
-XX:+ExitOnOutOfMemoryError

5. 配置数据源

  • presto支持不同的数据源,通过catalog进行配置;不同的数据源,有不同的catalog
  • 现以hive数据源为例,创建个hive的catalog
  • etc中创建目录catalog
cd /kkb/install/presto-server-0.237/etc
mkdir catalog
cd catalog
vim hive.properties
  • 添加如下内容
connector.name=hive-hadoop2
hive.metastore.uri=thrift://node03:9083

6. 分发presto

cd /kkb/install/
scp -r presto node02:/kkb/install/
scp -r presto node03:/kkb/install/

7. 配置node.properties

  • 进入三台节点的/kkb/install/presto/etc目录,修改node.properties文件
cd /kkb/install/presto/etc
vim node.properties
  • 三台节点的内容分别如下
# node01如下内容
node.environment=production
node.id=ffffffff-ffff-ffff-ffff-fffffffffff1
node.data-dir=/kkb/install/presto/data

# node2如下内容
node.environment=production
node.id=ffffffff-ffff-ffff-ffff-fffffffffff2
node.data-dir=/kkb/install/presto/data

# node03如下内容
node.environment=production
node.id=ffffffff-ffff-ffff-ffff-fffffffffff3
node.data-dir=/kkb/install/presto/data

说明:

node.environment 环境的名称;presto集群各节点的此名称必须保持一致

node.id presto每个节点的id,必须唯一

node.data-dir 存储log及其他数据的目录

8. 配置config.properties

  • 通过配置config.properties文件,指明server是coordinator还是worker

  • 虽然presto server可以同时作为coordinator和worker;但是为了更好的性能,一般让server要么作为coordinator,要么作为worker

  • presto是主从架构;主是coordinator,从是worker

  • 现设置node01作为coordinator节点;node02、node03节点作为worker节点

  • node01上配置coordinator

cd /kkb/install/presto/etc
vim config.properties
  • 添加如下内容
coordinator=true
node-scheduler.include-coordinator=false
http-server.http.port=8880
query.max-memory=50GB
query.max-memory-per-node=1GB
discovery-server.enabled=true
discovery.uri=http://node01:8880

说明:

coordinator=true 允许此presto实例作为coordinator

node-scheduler.include-coordinator 是否允许在coordinator上运行work

http-server.http.port presto使用http服务进行内部、外部的通信;指定http server的端口

query.max-memory 一个查询运行时,使用的所有的分布式内存的总量的上限

query.max-memory-per-node query在执行时,使用的任何一个presto服务器上使用的内存上限

discovery-server.enabled presto使用discovery服务,用来发现所有的presto节点

discovery.uri discovery服务的uri

  • node02、node03上配置worker
cd /kkb/install/presto/etc
vim config.properties
  • 添加如下内容
coordinator=false
http-server.http.port=8880
query.max-memory=50GB
discovery.uri=http://node01:8880

9. 启动presto server

  • 若要用presto对接hive数据,需要启动hive metastore服务
  • 上课环境:hive安装在node03上,所以在node03启动metastore服务
nohup hive --service metastore > /dev/null 2>&1 &
  • 在node01、node02、node03上分别启动presto server,执行以下命令
cd /kkb/install/presto
# 前台启动,控制台打印日志
bin/launcher run
# 或使用后台启动presto
bin/launcher start
  • jps查看,各节点出现名为PrestoServer的进程
  • 日志所在目录
/kkb/install/presto/data/var/log

2. 安装部署Presto命令行接口

1. 下载安装包

2. 重命名文件

cd /kkb/soft
mv presto-cli-0.237-executable.jar prestocli

3. 增加可执行权限

chmod u+x prestocli

4. 启动presto cli

  • 注意:先启动HDFS

  • 查看presto客户端jar包的使用方式

./prestocli --help
  • 两种方式;方式一
./prestocli --server node01:8880 --catalog hive --schema default

说明:

--catalog hive 中的hive指的是presto/etc/catalog中的hive.properties的文件名

  • 方式二
java -jar presto-cli-0.237-executable.jar --server node01:8880 --catalog hive --schema default
  • 退出presto cli
quit

5. 体验命令操作

Presto的命令行操作,相当于Hive命令行操作。每个表必须要加上schema前缀;例如

select * from schema.<tablename> limit 5;

或者切换到指定的schema,再查询表数据

use myhive;
select * from score limit 3;

异常处理:在node01进行presto查询时没有响应

presto:default> show tables;

Query 20191008_100336_00007_kdyy2, WAITING_FOR_RESOURCES, 0 nodes, 0 splits

解决办法:这是因为在前面的配置中node01只配置了coordinator,不允许worker同时运行,如果想要node01既运行coordinator又同时运行worker,需要修改配置文件:

vim ./etc/config.properties

...
node-scheduler.include-coordinator=true
...

3. 安装部署Presto 可视化客户端

1. 下载安装包

  • presto有个开源的带可视化界面的客户端yanagishima
  • 源码下载地址:yanagishima
  • 官网地址
  • 将下载的包yanagishima-18.0.zip上传到node01点/kkb/soft目录

2. 解压缩

cd /kkb/soft
unzip -d /kkb/install yanagishima-18.0.zip

# 若出现-bash: unzip: command not found,表示没有安装unzip;需要安装;然后再解压缩
sudo yum -y install unzip zip

cd /kkb/install/yanagishima-18.0

3. 修改配置文件

  • 修改yanagishima.properties文件
cd /kkb/install/yanagishima-18.0/conf
vim yanagishima.properties
  • 添加如下内容
jetty.port=7080
presto.datasources=kkb-presto
presto.coordinator.server.kkb-presto=http://node01:8880
catalog.kkb-presto=hive
schema.kkb-presto=default
sql.query.engines=presto

注意这里kkb-presto是数据源的名称,可以自定义

4. 启动yanagishima

nohup bin/yanagishima-start.sh >yanagishima.log 2>&1 &
  • 注意一定要进入yanagishima的安装目录下以相对路径方式执行以上命令启动

  • node01上多出名为YanagishimaServer的进程

  • 启动web界面

    http://node01:7080

    在界面中进行查询了

    若ui界面显示很慢,或者不显示,可以尝试将node01替换成相应的ip地址

  • 查看表结构;

  • 每个表后面都有个复制键,点一下会复制完整的表名,然后再上面框里面输入sql语句,ctrl+enter组合键或Run按钮执行显示结果

image-20200630111827193

  • 这里有个Tree View,可以查看所有表的结构,包括Schema、表、字段等。

  • 比如执行select * from hive.myhive.score,这个句子里Hive这个词可以删掉,即变成select * from myhive.score;hive是上面配置的Catalog名称

  • 特别注意:sql语句末尾千万不要加分号;否则报错

image-20200630112056143

5. Presto查询及优化

1. Presto sql语法

以下用presto的客户端连接hive connector来进行演示

如果使用UI界面演示去掉SQL语句最后面的分号即可

查看schema有哪些(对应hive中的database)

SHOW SCHEMAS;

进入schema

use myhive;

查看有哪些表

SHOW TABLES;

创建schema

语法:CREATE SCHEMA [ IF NOT EXISTS ] schema_name

CREATE SCHEMA testschema;

删除schema

语法:DROP SCHEMA [ IF EXISTS ] schema_name
drop schema testschema;

创建表

语法:CREATE TABLE [ IF NOT EXISTS ]
table_name (column_name data_type [ COMMENT comment],... ]

create table stu4(id int, name varchar(20));

创建表CTAS

语法:
CREATE TABLE [ IF NOT EXISTS ] table_name [ ( column_alias, ... ) ]
[ COMMENT table_comment ]
[ WITH ( property_name = expression [, ...] ) ]
AS query
[ WITH [ NO ] DATA ]

create table if not exists myhive.stu5 as select id, name from stu1;

删除表中符合条件的行

语法:DELETE FROM table_name [ WHERE condition ]
说明:hive connector只支持一次性的删除一个完整的分区;不支持删除一行数据

DELETE FROM order_partition where month='2019-03';

查看表的描述信息

DESCRIBE hive.myhive.stu1;

ANALYZE获得表及列的统计信息

语法:ANALYZE table_name

ANALYZE hive.myhive.stu1;

prepare 给statement起一个名称,等待将来的执行

execute执行一个准备好的statement

语法:PREPARE statement_name FROM statement

prepare my_select1 from select * from score;
execute my_select1;

prepare my_select2 from select * from score where s_score < 90 and s_score > 70;
execute my_select2;

prepare my_select3 from select * from score where s_score < ? and s_score > ?;
execute my_select3 using 90, 70;

EXPLAIN:查询一个statement的逻辑计划或分布式执行计划,或校验statement

语法:
EXPLAIN [ ( option [, ...] ) ] statement

where option can be one of:

    FORMAT { TEXT | GRAPHVIZ | JSON }
    TYPE { LOGICAL | DISTRIBUTED | VALIDATE | IO }

查询逻辑计划语句:
explain select s_id, avg(s_score) from score group by s_id;
等价于
explain (type logical)select s_id, avg(s_score) from score group by s_id;

查询分布式执行计划distributed execution plan
explain (type distributed)select s_id, avg(s_score) from score group by s_id;

校验语句的正确性
explain (type validate)select s_id, avg(s_score) from score group by s_id;

explain (type io, format json)select s_id, avg(s_score) from score group by s_id;

SELECT查询

语法:
[ WITH with_query [, ...] ]
SELECT [ ALL | DISTINCT ] select_expr [, ...]
[ FROM from_item [, ...] ]
[ WHERE condition ]
[ GROUP BY [ ALL | DISTINCT ] grouping_element [, ...] ]
[ HAVING condition]
[ { UNION | INTERSECT | EXCEPT } [ ALL | DISTINCT ] select ]
[ ORDER BY expression [ ASC | DESC ] [, ...] ]
[ LIMIT [ count | ALL ] ]

from_item:
table_name [ [ AS ] alias [ ( column_alias [, ...] ) ] ]
from_item join_type from_item [ ON join_condition | USING ( join_column [, ...] ) ]

join_type:
[ INNER ] JOIN
LEFT [ OUTER ] JOIN
RIGHT [ OUTER ] JOIN
FULL [ OUTER ] JOIN
CROSS JOIN

grouping_element:
()
expression
GROUPING SETS ( ( column [, ...] ) [, ...] )
CUBE ( column [, ...] )
ROLLUP ( column [, ...] )

语句:
with语句:用于简化内嵌的子查询
select a, b
from (
select s_id as a, avg(s_score) as b from score group by s_id
) as tbl1;

等价于:
with tbl1 as (select s_id as a, avg(s_score) as b from score group by s_id)
select a, b from tbl1;

多个子查询也可以用with
WITH
  t1 AS (SELECT a, MAX(b) AS b FROM x GROUP BY a),
  t2 AS (SELECT a, AVG(d) AS d FROM y GROUP BY a)
SELECT t1.*, t2.*
FROM t1
JOIN t2 ON t1.a = t2.a;

with语句中的关系可以串起来(chain)
WITH
  x AS (SELECT a FROM t),
  y AS (SELECT a AS b FROM x),
  z AS (SELECT b AS c FROM y)
SELECT c FROM z;

group by:
select s_id as a, avg(s_score) as b from score group by s_id;
等价于:
select s_id as a, avg(s_score) as b from score group by 1;
1代表查询输出中的第一列s_id

select count(*) as b from score group by s_id;

2. 存储优化

  • 合理设置分区

    与Hive类似,Presto会根据元信息读取分区数据,合理的分区能减少Presto数据读取量,提升查询性能。

  • 使用列式存储

    Presto对ORC文件读取做了特定优化,因此在Hive中创建Presto使用的表时,建议采用ORC格式存储。相对于Parquet,Presto对ORC支持更好。

  • 使用压缩

    数据压缩可以减少节点间数据传输对IO带宽压力,对于即席查询需要快速解压,建议采用snappy压缩

  • 预先排序

    对于已经排序的数据,在查询的数据过滤阶段,ORC格式支持跳过读取不必要的数据。比如对于经常需要过滤的字段可以预先排序。

3. SQL优化

  • 列剪裁

    只选择使用必要的字段: 由于采用列式存储,选择需要的字段可加快字段的读取、减少数据量。避免采用*读取所有字段

[GOOD]: SELECT s_id, c_id FROM score

[BAD]:  SELECT * FROM score
  • 过滤条件必须加上分区字段

    对于分区表,where语句中优先使用分区字段进行过滤。day是分区字段,vtime是具体访问时间

[GOOD]: SELECT vtime, stu, address FROM tbl where day=20200501

[BAD]:  SELECT * FROM tbl where vtime=20200501
  • Group By语句优化:

    合理安排Group by语句中字段顺序对性能有一定提升。将Group By语句中字段按照每个字段distinct数据多少进行降序排列, 减少GROUP BY语句后面的排序一句字段的数量能减少内存的使用.

uid个数多;gender少
[GOOD]: SELECT GROUP BY uid, gender

[BAD]:  SELECT GROUP BY gender, uid
  • Order by时使用Limit, 尽量避免ORDER BY: Order by需要扫描数据到单个worker节点进行排序,导致单个worker需要大量内存
[GOOD]: SELECT * FROM tbl ORDER BY time LIMIT 100

[BAD]:  SELECT * FROM tbl ORDER BY time
  • 使用近似聚合函数: 对于允许有少量误差的查询场景,使用这些函数对查询性能有大幅提升。比如使用approx_distinct() 函数比Count(distinct x)有大概2.3%的误差
select approx_distinct(s_id) from score;
  • 用regexp_like代替多个like语句: Presto查询优化器没有对多个like语句进行优化,使用regexp_like对性能有较大提升
SELECT
...
FROM
access
WHERE
method LIKE '%GET%' OR
method LIKE '%POST%' OR
method LIKE '%PUT%' OR
method LIKE '%DELETE%'

优化:
SELECT
...
FROM
access
WHERE
regexp_like(method, 'GET|POST|PUT|DELETE')
  • 使用Join语句时将大表放在左边: Presto中join的默认算法是broadcast join,即将join左边的表分割到多个worker,然后将join右边的表数据整个复制一份发送到每个worker进行计算。如果右边的表数据量太大,则可能会报内存溢出错误。
[GOOD] SELECT ... FROM large_table l join small_table s on l.id = s.id
[BAD] SELECT ... FROM small_table s join large_table l on l.id = s.id
  • 使用Rank函数代替row_number函数来获取Top N

  • UNION ALL 代替 UNION :不用去重

  • 使用WITH语句: 查询语句非常复杂或者有多层嵌套的子查询,请试着用WITH语句将子查询分离出来

6. 其他注意事项

1. 字段名引用

  • 避免和关键字冲突:MySQL对字段加反引号`;Presto对字段加双引号分割

    当然,如果字段名称不是关键字,可以不加这个双引号。

2. 函数

  • 对于Timestamp,需要进行比较的时候,需要添加Timestamp关键字,而MySQL中对Timestamp可以直接进行比较。
/*MySQL的写法*/
SELECT t FROM a WHERE t > '2020-05-01 00:00:00'; 

/*Presto的写法*/
SELECT t FROM a WHERE t > timestamp '2020-05-01 00:00:00';

3. 不支持INSERT OVERWRITE语法

  • Presto中不支持insert overwrite语法,只能先delete,然后insert into。

4. QUET格式

  • Presto目前支持Parquet格式,支持查询,但不支持insert

五、拓展

Views: 28

Maxwell 数据库数据实时采集

1、Maxwell 简介

Maxwell 是一个能实时读取 MySQL 二进制日志文件binlog,并生成 Json格式的消息,作为生产者发送给 Kafka,Kinesis、RabbitMQ、Redis、Google Cloud Pub/Sub、文件或其它平台的应用程序。它的常见应用场景有ETL、维护缓存、收集表级别的dml指标、增量到搜索引擎、数据分区迁移、切库binlog回滚方案等。

Maxwell主要提供了下列功能

    1. 支持SELECT * FROM table的方式进行全量数据初始化。
    1. 支持在主库发生failover后,自动恢复binlog位置,实现断点续传。
    1. 可以对数据进行分区,解决数据倾斜问题,发送到Kafka的数据支持库、表、列等级别的数据分区。
    1. 工作方式是伪装为slave接收binlog events,然后根据schema信息拼装,可以接受ddl、xid、row等event。

2、Mysql Binlog介绍

2.1 Binlog 简介

MySQL中一般有以下几种日志

日志类型 写入日志的信息
错误日志 记录在启动,运行或停止mysqld时遇到的问题
通用查询日志 记录建立的客户端连接和执行的语句
二进制日志 binlog 记录更改数据的语句
中继日志 从服务器 复制 主服务器接收的数据更改
慢查询日志 记录所有执行时间超过 long_query_time 秒的所有查询或不使用索引的查询
DDL日志(元数据日志) 元数据操作由DDL语句执行

在默认情况下,系统仅仅打开错误日志,关闭了其他所有日志,以达到尽可能减少IO损耗提高系统性能的目的,但是在一般稍微重要一点的实际应用场景中,都至少需要打开二进制日志,因为这是MySQL很多存储引擎进行增量备份的基础,也是MySQL实现复制的基本条件

接下来主要介绍二进制日志 binlog。

MySQL 的二进制日志 binlog 可以说是 MySQL 最重要的日志,它记录了所有的 DDLDML 语句(除了数据查询语句select、show等),以事件形式记录,还包含语句所执行的消耗的时间,MySQL的二进制日志是事务安全型的。binlog 的主要目的是复制和恢复

Binlog日志的两个最重要的使用场景

  • MySQL主从复制
    • MySQL Replication在Master端开启binlog,Master把它的二进制日志传递给slaves来达到master-slave数据一致的目的。
  • 数据恢复
    • 通过使用 mysqlbinlog工具来使恢复数据。

2.2 Binlog 的日志格式

记录在二进制日志中的事件的格式取决于二进制记录格式。支持三种格式类型:

  • Statement:基于SQL语句的复制(statement-based replication, SBR)
  • Row:基于行的复制(row-based replication, RBR)
  • Mixed:混合模式复制(mixed-based replication, MBR)

Statement

  • 每一条会修改数据的sql都会记录在binlog中。
  • 优点
    • 不需要记录每一行的变化,减少了binlog日志量,节约了IO, 提高了性能。
  • 缺点
    • 在进行数据同步的过程中有可能出现数据不一致。
    • 比如 update tt set create_date=now(),如果用binlog日志进行恢复,由于执行时间不同可能产生的数据就不同。

Row

  • 它不记录sql语句上下文相关信息,仅保存哪条记录被修改。
  • 优点
    • 保持数据的绝对一致性。因为不管sql是什么,引用了什么函数,它只记录执行后的效果。
  • 缺点
    • 每行数据的修改都会记录,最明显的就是update语句,导致更新多少条数据就会产生多少事件,占用较大空间。

Mixed

  • 从5.1.8版本开始,MySQL提供了Mixed格式,实际上就是Statement与Row的结合。
  • 在Mixed模式下,一般的复制使用Statement模式保存binlog,对于Statement模式无法复制的操作使用Row模式保存binlog, MySQL会根据执行的SQL语句选择日志保存方式(因为statement只有sql,没有数据,无法获取原始的变更日志,所以一般建议为Row模式)。
  • 优点
    • 节省空间,同时兼顾了一定的一致性。
  • 缺点
    • 还有些极个别情况依旧会造成不一致,另外statement和mixed对于需要对binlog的监控的情况都不方便。

3、Mysql 实时数据同步方案对比

  • mysql 数据实时同步可以通过解析mysql的 binlog 的方式来实现,解析binlog可以有多种方式,可以通过canal,或者maxwell等各种方式实现。以下是各种抽取方式的对比介绍。

    mysql实时同步方案对比

  • 其中canal 由 Java开发,分为服务端和客户端,拥有众多的衍生应用,性能稳定,功能强大;canal 需要自己编写客户端来消费canal解析到的数据。

  • Maxwell相对于canal的优势是使用简单,Maxwell比Canal更加轻量级,它直接将数据变更输出为json字符串,不需要再编写客户端。对于缺乏基础建设,短时间内需要快速迭代的项目和公司比较合适。

  • 另外Maxwell 有一个亮点功能,就是Canal只能抓取最新数据,对已存在的历史数据没有办法处理。而Maxwell有一个bootstrap功能,可以直接引导出完整的历史数据用于初始化,非常好用。

4、开启Mysql的Binlog

  • 1、服务器当中安装mysql(省略)

    • 注意:mysql的版本尽量不要太低,也不要太高,最好使用5.6及以上版本。
  • 2、添加mysql普通用户maxwell

    • 为mysql添加一个普通用户maxwell,因为maxwell这个软件默认用户使用的是maxwell这个用户。

    • 进入mysql客户端,然后执行以下命令,进行授权

    mysql -uroot -p123456
    • 执行sql语句
    --校验级别最低,只校验密码长度
    mysql> set global validate_password_policy=LOW;
    mysql> set global validate_password_length=6;
    
    --创建maxwell库(启动时候会自动创建,不需手动创建)和用户
    mysql> CREATE USER 'maxwell'@'%' IDENTIFIED BY '123456';
    mysql> GRANT ALL ON maxwell.* TO 'maxwell'@'%';
    mysql> GRANT SELECT, REPLICATION CLIENT, REPLICATION SLAVE on *.* to 'maxwell'@'%'; 
    --刷新权限
    mysql> flush privileges;

    maxwell会自动在MySQL中创建名为maxwell的数据库作为元数据保存使用。

  • 3、修改配置文件 /etc/my.cnf

    • 执行命令 sudo vim /etc/my.cnf, 添加或修改以下三行配置
    #binlog日志名称前缀
    log-bin= /var/lib/mysql/mysql-bin
    
    #binlog日志格式
    binlog-format=ROW
    
    #唯一标识,这个值的区间是:1到(2^32)-1
    server_id=1
  • 4、重启mysql服务

    • 执行如下命令
    sudo service mysqld restart
  • 5、验证binlog是否配置成功

    • 进入mysql客户端,并执行以下命令进行验证
    mysql -uroot -p123456
    mysql> show variables like '%log_bin%';

    image-20210517161330004

  • 6、查看binlog日志文件生成

    • 进入 /var/lib/mysql 目录,查看binlog日志文件.

    image-20210517162315281

5、Maxwell安装部署

  • 1、下载对应版本的安装包

  • 2、上传服务器

  • 3、解压安装包到指定目录

    tar -zxvf maxwell-1.21.1.tar.gz -C /kkb/install/
  • 4、修改maxwell配置文件

    • 进入到安装目录 /kkb/install/maxwell-1.21.1 进行如下操作
    cd /kkb/install/maxwell-1.21.1 
    cp config.properties.example config.properties
    vim config.properties
    • 配置文件config.properties 内容如下:
    # choose where to produce data to
    producer=kafka
    # list of kafka brokers
    kafka.bootstrap.servers=node01:9092,node02:9092,node03:9092
    # mysql login info
    host=node03
    port=3306
    user=maxwell
    password=123456
    # kafka topic to write to
    kafka_topic=maxwell
    • 注意:一定要保证使用maxwell 用户和 123456 密码能够连接上mysql数据库。

6、kafka介绍和使用

6.1 Kafka简介

​ Kafka是最初由Linkedin公司开发,它是一个分布式、可分区、多副本,基于zookeeper协调的分布式日志系统;常见可以用于web/nginx日志、访问日志,消息服务等等。Linkedin于2010年贡献给了Apache基金会并成为顶级开源项目。主要应用场景是:日志收集系统和消息系统

​ Kafka是一个分布式消息队列。具有高性能、持久化、多副本备份、横向扩展能力。生产者往队列里写消息,消费者从队列里取消息进行业务逻辑。Kafka就是一种发布-订阅模式。将消息保存在磁盘中,以顺序读写方式访问磁盘,避免随机读写导致性能瓶颈。

  • 消息(Message)
    • 是指在应用之间传送的数据,消息可以非常简单,比如只包含文本字符串,也可以更复杂,可能包含嵌入对象。
  • 消息队列(Message Queue)
    • 一种应用间的通信方式,消息发送后可以立即返回,通过消息系统来确保信息的可靠传递,消息发布者只管把消息发布到MQ中而不管谁来取,消息使用者只管从MQ中取消息而不管谁发布的,这样发布者和使用者都不用知道对方的存在。

6.2 Kafka特性

  • 高吞吐、低延迟

    kafka 最大的特点就是收发消息非常快,kafka 每秒可以处理几十万条消息,它的最低延迟只有几毫秒。
  • 高伸缩性

    每个主题(topic) 包含多个分区(partition),主题中的分区可以分布在不同的主机(broker)中。
  • 持久性、可靠性

    Kafka 能够允许数据的持久化存储,消息被持久化到磁盘,并支持数据备份防止数据丢失。
  • 容错性

    允许集群中的节点失败,某个节点宕机,Kafka 集群能够正常工作。
  • 高并发

    支持数千个客户端同时读写。

6.3 Kafka集群架构

kafka集群架构

  • producer

    消息生产者,发布消息到Kafka集群的终端或服务。
  • broker

    Kafka集群中包含的服务器,一个borker就表示kafka集群中的一个节点。
  • topic

    每条发布到Kafka集群的消息属于的类别,即Kafka是面向 topic 的。
    更通俗的说Topic就像一个消息队列,生产者可以向其写入消息,消费者可以从中读取消息,一个Topic支持多个生产者或消费者同时订阅它,所以其扩展性很好。
  • partition

    每个 topic 包含一个或多个partition。Kafka分配的单位是partition。
  • replica

    partition的副本,保障 partition 的高可用。
  • consumer

    从Kafka集群中消费消息的终端或服务。
  • consumer group

    每个 consumer 都属于一个 consumer group,每条消息只能被 consumer group 中的一个 Consumer 消费,但可以被多个 consumer group 消费。
  • leader

    每个partition有多个副本,其中有且仅有一个作为Leader,Leader是当前负责数据的读写的partition。 producer 和 consumer 只跟 leader 交互。
  • follower

    Follower跟随Leader,所有写请求都通过Leader路由,数据变更会广播给所有Follower,Follower与Leader保持数据同步。如果Leader失效,则从Follower中选举出一个新的Leader。
  • controller

    知道大家有没有思考过一个问题,就是Kafka集群中某个broker宕机之后,是谁负责感知到他的宕机,以及负责进行Leader Partition的选举?如果你在Kafka集群里新加入了一些机器,此时谁来负责把集群里的数据进行负载均衡的迁移?包括你的Kafka集群的各种元数据,比如说每台机器上有哪些partition,谁是leader,谁是follower,是谁来管理的?如果你要删除一个topic,那么背后的各种partition如何删除,是谁来控制?还有就是比如Kafka集群扩容加入一个新的broker,是谁负责监听这个broker的加入?如果某个broker崩溃了,是谁负责监听这个broker崩溃?这里就需要一个Kafka集群的总控组件,Controller。他负责管理整个Kafka集群范围内的各种东西。
    
  • zookeeper

    (1)   Kafka 通过 zookeeper 来存储集群的meta元数据信息。
    (2)一旦controller所在broker宕机了,此时临时节点消失,集群里其他broker会一直监听这个临时节点,发现临时节点消失了,就争抢再次创建临时节点,保证有一台新的broker会成为controller角色。
  • offset

    • 偏移量
    消费者在对应分区上已经消费的消息数(位置),offset保存的地方跟kafka版本有一定的关系。
    kafka0.8 版本之前offset保存在zookeeper上。
    kafka0.8 版本之后offset保存在kafka集群上。
    它是把消费者消费topic的位置通过kafka集群内部有一个默认的topic,
    名称叫 __consumer_offsets,它默认有50个分区。

6.4 Kafka集群安装部署

  • 1、下载安装包(http://kafka.apache.org

    kafka_2.11-1.1.0.tgz
  • 2、规划安装目录

    /kkb/install
  • 3、上传安装包到服务器中

    通过FTP工具上传安装包到node01服务器上
  • 4、解压安装包到指定规划目录

    tar -zxvf kafka_2.11-1.1.0.tgz -C /kkb/install
  • 5、重命名解压目录

    mv kafka_2.11-1.1.0 kafka
  • 6、修改配置文件

    • 在node01上修改

    • 进入到kafka安装目录下有一个config目录

      • vi server.properties
      #指定kafka对应的broker id ,唯一
      broker.id=0
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node01
    • 配置kafka环境变量

      • sudo vi /etc/profile
      export KAFKA_HOME=/kkb/install/kafka
      export PATH=$PATH:$KAFKA_HOME/bin
  • 6、分发kafka安装目录到其他节点

    scp -r kafka node02:/kkb/install
    scp -r kafka node03:/kkb/install
    scp /etc/profile node02:/etc
    scp /etc/profile node03:/etc
  • 7、修改node02和node03上的配置

    • node02

    • vi server.properties

      #指定kafka对应的broker id ,唯一
      broker.id=1
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node02
    • node03

    • vi server.properties

      #指定kafka对应的broker id ,唯一
      broker.id=2
      #指定数据存放的目录
      log.dirs=/kkb/install/kafka/kafka-logs
      #指定zk地址
      zookeeper.connect=node01:2181,node02:2181,node03:2181
      #指定是否可以删除topic ,默认是false 表示不可以删除
      delete.topic.enable=true
      #指定broker主机名
      host.name=node03
  • 8、让每台节点的kafka环境变量生效

    • 在每台服务器执行命令
    source /etc/profile

6.5 kafka集群启动和停止

  • 1、启动kafka集群

    • 先启动zookeeper集群,然后在所有节点如下执行脚本
    nohup kafka-server-start.sh /kkb/install/kafka/config/server.properties >/dev/null 2>&1 &
  • 2、停止kafka集群

    • 所有节点执行关闭kafka脚本
    kafka-server-stop.sh

6.6 kafka命令行的管理使用

  • 1、创建topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --create --partitions 3 --replication-factor 2 --topic test --zookeeper node01:2181,node02:2181,node03:2181
  • 2、查询所有的topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --list --zookeeper node01:2181,node02:2181,node03:2181 
  • 3、查看topic的描述信息

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --describe --topic test --zookeeper node01:2181,node02:2181,node03:2181  
  • 4、删除topic

    • 使用 kafka-topics.sh脚本
    kafka-topics.sh --delete --topic test --zookeeper node01:2181,node02:2181,node03:2181 
  • 5、模拟生产者写入数据到topic中

    • 使用 kafka-console-producer.sh 脚本
    kafka-console-producer.sh --broker-list node01:9092,node02:9092,node03:9092 --topic test 
  • 6、模拟消费者拉取topic中的数据

    • 使用 kafka-console-consumer.sh 脚本
    kafka-console-consumer.sh --zookeeper node01:2181,node02:2181,node03:2181 --topic test --from-beginning

    或者(推荐)

    kafka-console-consumer.sh --bootstrap-server node01:9092,node02:9092,node03:9092 --topic test --from-beginning

7、Maxwell实时采集mysql表数据到kafka

  • 1、启动kafka集群和zookeeper集群

    • 启动zookeeper集群
    #每台节点执行脚本
    nohup zkServer.sh start >/dev/null  2>&1 &
    • 启动kafka集群
    nohup /kkb/install/kafka/bin/kafka-server-start.sh /kkb/install/kafka/co
    nfig/server.properties > /dev/null 2>&1 &
  • 2、创建topic

    kafka-topics.sh --create --topic maxwell --partitions 3 --replication-factor 2 --zookeeper node01:2181,node02:2181,node03:2181

    (如果虚拟机磁盘容量有效,可以将分区数和复制因子都设置为1.)

  • 3、启动maxwell服务

    /kkb/install/maxwell-1.21.1/bin/maxwell
  • 4、插入数据并进行测试

    • 向mysql表中插入一条数据,并开启kafka的消费者,查看kafka是否能够接收到数据。

    • 向mysql当中创建数据库和数据库表并插入数据

      CREATE DATABASE /*!32312 IF NOT EXISTS*/<code>test_db /*!40100 DEFAULT CHARACTER SET utf8 */;
      
      USE test_db;
      
      /*Table structure for table user */
      
      DROP TABLE IF EXISTS user;
      
      CREATE TABLE user (
      id varchar(10) NOT NULL,
      name varchar(10) DEFAULT NULL,
      age int(11) DEFAULT NULL,
      PRIMARY KEY (id)
      ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
      
      /*Data for the table user */
      #插入数据
      insert  into user(id,name,age) values  ('1','xiaokai',20);
      #修改数据
      update user set age= 30 where id='1';
      #删除数据
      delete from user where id='1';
  • 5、启动kafka的自带控制台消费者

    • 启动Kafka消费者监听maxwell主题
    kafka-console-consumer.sh --topic maxwell --bootstrap-server node01:9092,node02:9092,node03:9092 --from-beginning 
    • 等待一段事件,观察maxwell主题是否有消息发送过来
    {"database":"test_db","table":"user","type":"insert","ts":1621244407,"xid":985,"commit":true,"data":{"id":"1","name":"xiaokai","age":20}}
    
    {"database":"test_db","table":"user","type":"update","ts":1621244413,"xid":999,"commit":true,"data":{"id":"1","name":"xiaokai","age":30},"old":{"age":20}}
    
    {"database":"test_db","table":"user","type":"delete","ts":1621244419,"xid":1013,"commit":true,"data":{"id":"1","name":"xiaokai","age":30}}
    • json数据字段说明

    • database

      • 数据库名称
    • table

      • 表名称
    • type

      • 操作类型
      • 包括 insert/update/delete 等
    • ts

      • 操作时间戳
    • xid

      • 事务id
    • commit

      • 同一个xid代表同一个事务,事务的最后一条语句会有commit
    • data

      • 最新的数据,修改后的数据
    • old

      • 旧数据,修改前的数据

Views: 30

Index