Kafka分布式消息系统

什么是Kafka

在Spark生态体系中,Kafka占有非常重要的位置。Kafka是一个使用Scala语言编写的基于ZooKeeper的高吞吐量低延迟的分布式发布与订阅消息系统,它可以实时处理大量消息数据以满足各种需求。比如基于Hadoop的批处理系统,低延迟的实时系统等。即便使用非常普通的硬件,Kafka每秒也可以处理数百万条消息,其延迟最低只有几毫秒。
在实际开发中,Kafka常常作为Spark Streaming的实时数据源,Spark Streaming从Kafka中读取实时消息进行处理,保证了数据的可靠性与实时性。二者是实时消息处理系统的重要组成部分。
那么Kafka到底是什么?简单来说,Kafka是消息中间件的一种。

举个生产者与消费者的例子:生产者生产鸡蛋,消费者消费鸡蛋。

假设消费者消费鸡蛋的时候噎住了(系统宕机了),而生产者还在生产鸡蛋,那么新生产的鸡蛋就丢失了;

再比如,生产者1秒钟生产100个鸡蛋(大交易量的情况),而消费者1秒钟只能消费50个鸡蛋,那过不了多长时间,消费者就吃不消了(消息堵塞,最终导致系统超时),导致鸡蛋又丢失了。

这个时候我们放个篮子在生产者与消费者中间,生产者生产出来的鸡蛋都放到篮子里,消费者去篮子里拿鸡蛋,这样鸡蛋就不会丢失了,这个篮子就相当于“Kafka”;鸡蛋则相当于Kafka中的消息(Message);篮子相当于存放消息的消息队列,也就是Kafka集群;

当篮子满了,鸡蛋放不下了,这时再加几个篮子,就是Kafka集群扩容。

Kafka中的一些基本概念:

消息(Message)

Kafka的数据单元被称为消息, 也称为事件(Event)。可以把消息看成是数据库里的一行数据或一条记录。为了提高效率,消息可以分组传输,每一组消息就是一个批次,分成批次传输可以减少网络开销。

服务器节点(Broker)

Kafka集群包含一个或多个服务器节点,一个独立的服务器节点被称为Broker。

主题(Topic)

每条发布到Kafka集群的消息都有一个类别,这个类别被称为主题。在物理上,不同主题的消息分开存储;在逻辑上,一个主题的消息虽然保存于一个或多个Broker上,但用户只需指定消息的主题即可生产或消费消息而不必关心消息存于何处。

分区(Partition)

为了使Kafka的吞吐率可以水平扩展,物理上把主题分成一个或多个分区。创建主题时可指定分区数量。

生产者(Producer)

负责发布消息到Kafka的Broker,实际上属于Broker的一种客户端。生产者负责选择哪些消息应该分配到哪个主题内的哪个分区。默认生产者会把消息均匀的分布到特定主题的所有分区上,但在某些情况下,生产者会将消息直接写到指定的分区。

消费者(Consumer)

从Kafka的Broker上读取消息的客户端。读取消息时需要指定读取的主题,通常消费者会订阅一个或多个主题,并按照消息生成的顺序读取他们。

不同的主题好比不同的高速公路,分区好比某条高速公路上的车道,消息就是车道上运行的车辆。如果车流量大,则拓宽车道,反之,则减少车道;而消费者就好比高速公路上的收费站,开放的收费站越多,则车辆通过速度越快。

Kafka架构

Kafka的消息传递流程如图所示。生产者将消息发送给Kafka集群,同时Kafka集群将消息转发给消费者。
一个典型的Kafka集群中包含若干生产者(数据可以是Web前端产生的页面内容或者服务器日志等)、若干Broker、若干消费者(可以是Hadoop集群、实时监控程序、数据仓库或其它服务)以及一个ZooKeeper集群。ZooKeeper用于管理和协调Broker。当Kafka系统中新增了Broker或者某个Broker故障失效时,ZooKeeper将通知生产者和消费者。生产者和消费者据此开始与其它Broker协调工作。生产者使用Push模式将消息发送到Broker,而消费者使用Pull模式从Broker订阅并消费消息。

file

file

主题与分区

Kafka通过主题对消息进行分类,一个主题可以分为多个分区,且每个分区可以存储于不同的Broker上,也就是说,一个主题可以横跨多个服务器。

对主题进行分区的好处是:允许主题消息规模超出一台服务器的文件大小上限。因为一个主题可以有多个分区,且可以存储在不同的服务器上,当一个分区的文件大小超出了所在服务器的文件大小上限时,可以动态添加其它分区,因此可以处理无限量的数据。

file

Kafka会为每个主题维护一个分区日志,记录各个分区的消息存放情况。消息以追加的方式写入到每个分区的尾部,然后以先入先出的顺序进行读取。由于一个主题包含多个分区,所以无法在整个主题范围内保证消息的顺序,但可以保证单个分区内消息的顺序。

当一条消息被发送到Broker时,会根据分区规则被存储到某个分区里。如果分区规则设置的合理,所有消息将被均匀的分配到不同的分区里,这样就实现了水平扩展。如果一个主题的消息都存放到一个文件中,则该文件所在的Broker的I/O将成为主题的性能瓶颈,而分区正好解决了这个问题。

分区中的每个记录都被分配了一个偏移量(offset),偏移量是一个连续递增的整数值,它唯一标识分区中的某个记录。而消费者只需保存该偏移量即可,当消费者客户端向Broker发起消息请求时需要携带偏移量。例如,消费者向Broker请求主题test的分区0中的偏移量从20开始的所有消息以及主题test的分区1中的偏移量从35开始的所有消。当消费者读取消息后,偏移量会线性递增。当然,消费者也可以按照任意顺序消费消息,比如读取已经消费过的历史消息(将偏移量重置到之前版本)。此外,消费者还可以指定从某个分区中一次最多返回多少条数据,防止一次返回数据太多而耗尽客户端的内存。

file

分区副本

在Kafka集群中,为了提高数据的可靠性,同一个分区可以复制多个副本分配到不同的Broker,这种方式类似于HDFS中的副本机制。如果其中一个Broker宕机,其它Broker可以接替宕机的Broker,不过生产者和消费者需要重新连接到新的Broker。

Kafka每个分区的副本都被分为两种类型:领导者副本和跟随者副本。领导者副本只有一个,其余的都是跟随者副本。所有生产者和消费者都向领导者副本发起请求,进行消息的写入与读取,而跟随者副本并不处理客户端的请求,它唯一的任务是从领导者副本复制消息,以保持与领导者副本数据及状态的一致。
如果领导者副本发生崩溃,会从其余的跟随者副本中选出一个作为新的领导者副本。

file

file

消费者组

消费者组(Consumer Group)实际上就是一组消费者的集合。每个消费者属于一个特定的消费者组(可为每个消费者指定组名称,消费者通过组名称对自己进行标识,若不指定组名称则属于默认的组)。
传统消息处理有两种模式:队列模式和发布订阅模式。队列模式是指消费者可以从一台服务器读取消息,并且每个消息只被其中一个消费者消费;发布订阅模式是指消息通过广播方式发送给所有消费者。而Kafka提供了消费者组模式,能够同时具备这两种(队列和发布订阅)模式的特点。
Kafka规定,同一消费者组内不允许多个消费者消费同一分区的消息;而不同的消费者组,可以同时消费同一分区的消息。也就是说,分区与同一个消费者组中的消费者的对应关系是多对一而不允许一对多。举个例子,如果同一个应用有100台机器,这100台机器属于同一个消费者组,则同一条消息在100台机器中只有一台能得到。如果另一个应用也需要同时消费同一个主题的消息,则需要新建一个消费者组并消费同一个主题的消息。我们已经知道,消息存储于分区中,消费者组与分区的关系如图。

file

Kafka集群环境搭建

Kafka依赖ZooKeeper集群,搭建Kafka集群之前,需要先搭建好ZooKeeper集群。ZooKeeper集群的搭建步骤此处不做过多讲解。本例依然使用三台服务器在CentOS7上搭建Kafka集群,三台服务器的主机和IP地址分别为:

centos01 192.168.170.133
centos02 192.168.170.134
centos03 192.168.170.135

1、下载解压Kafka

从Apache官网http://kafka.apache.org下载Kafka的稳定版本kafka_2.11-2.0.0.tgz。
然后将Kafka安装包上传到centos01节点的/opt/softwares目录,并解压到目录/opt/modules下:

$ tar -zxvf kafka_2.11-2.0.0.tgz -C /opt/modules/

2、修改配置文件

修改Kafka安装目录下的config/server.properties文件。在分布式环境中建议至少修改以下配置项:

broker.id=1
num.partitions=2
default.replication.factor=2
listeners=PLAINTEXT://centos01:9092
log.dirs=/opt/modules/kafka_2.11-2.0.0/kafka-logs
zookeeper.connect=centos01:2181,centos02:2181,centos03:2181

3、发送安装文件到其他节点
将centos01节点配置好的Kafka安装文件复制到centos02和centos03节点:

scp -r kafka_2.11-2.0.0/ hadoop@centos02:/opt/modules/
scp -r kafka_2.11-2.0.0/ hadoop@centos03:/opt/modules/

复制完成后,修改centos02节点的Kafka安装目录下的config/server.properties文件,修改内容如下:

broker.id=2
listeners=PLAINTEXT://centos02:9092

同理,修改centos03节点的Kafka安装目录下的config/server.properties文件,修改内容如下:

broker.id=3
listeners=PLAINTEXT://centos03:9092

4、启动ZooKeeper集群
分别在三个节点上执行以下命令,启动ZooKeeper集群(需进入ZooKeeper安装目录):

bin/zkServer.sh start

5、启动Kafka集群
分别在三个节点上执行以下命令,启动Kafka集群(需进入Kafka安装目录):

bin/kafka-server-start.sh -daemon config/server.properties

集群启动后,分别在各个节点上执行jps命令,查看启动的Java进程,若能输出如下进程信息,说明启动成功。

2848 Jps
2518 QuorumPeerMain
2795 Kafka

查看Kafka安装目录下的日志文件logs/server.log,确保运行稳定,没有抛出异常。至此,Kafka集群搭建完成。

Kafka命令行操作

1、创建主题
创建主题可以使用Kafka提供的命令工具kafka-topics.sh,此处我们创建一个名为topictest的主题,分区数为2,每个分区的副本数为2,命令如下(在Kafka集群的任意节点执行即可):

$ bin/kafka-topics.sh \
--create \
--zookeeper centos01:2181,centos02:2181,centos03:2181 \
--replication-factor 2 \
--partitions 2 \
--topic topictest

说明

  • --create:指定命令的动作是创建主题,使用该命令必须指定--topic参数。
  • --topic:所创建的主题名称。
  • --partitions:所创建主题的分区数。
  • --zookeeper:指定ZooKeeper集群的访问地址,这种方式已废弃
    • 建议改成--bootstrap-server的方式指定broker所在地址和端口,如有多个用逗号分隔。
  • --replication-factor:所创建主题的分区副本数,其值必须小于等于Kafka的节点数。
    • 如果只有一个节点,可以不指定

命令执行完毕后,若输出以下结果则表明创建主题成功:

Created Topic "topictest".

2、查询主题
创建主题成功后,可以执行以下命令,查看当前Kafka集群中存在的所有主题:

$ bin/kafka-topics.sh \
--list \
--zookeeper centos01:2181

也可以使用--describe参数查询某一个主题的详细信息。例如,查询主题topictest的详细信息,命令如下:

$ bin/kafka-topics.sh \
--describe \
--zookeeper centos01:2181 \
--topic topictest

输出结果如下:

Topic:topictest PartitionCount:2    ReplicationFactor:2 Configs:
Topic: topictest    Partition: 0    Leader: 2   Replicas: 2,3   Isr: 2,3
Topic: topictest    Partition: 1    Leader: 3   Replicas: 3,1   Isr: 3,1

可以看到,该主题有2个分区,每个分区有2个副本。分区编号为0的副本分布在broker.id为2和3的Broker上,其中broker.id为2上的副本为领导者副本;分区编号为1的副本分布在broker.id为1和3的Broker上,其中broker.id为3上的副本为领导者副本。

3、创建生产者
Kafka生产者作为消息生产角色,可以使用Kafka自带的命令工具创建一个最简单的生产者。例如,在主题topictest上创建一个生产者,命令如下:

$ bin/kafka-console-producer.sh \
--broker-list centos01:9092,centos02:9092,centos03:9092 \
--topic topictest

说明:

  • --broker-list:指定Kafka Broker的访问地址,只要能访问到其中一个即可连接成功,若想写多个则用逗号隔开。建议将所有的Broker都写上,如果只写其中一个,如果该Broker失效,连接将失败。注意此处的Broker访问端口为9092,Broker通过该端口接收生产者和消费者的请求,该端口在安装Kafka时已经指定。
  • --topic:指定生产者发送消息的主题名称。
    创建完成后,控制台进入等待键盘输入消息的状态。

4、创建消费者
新开启一个SSH连接窗口(可连接Kafka集群中的任何一个节点),在主题topictest上创建一个消费者,命令如下:

$ bin/kafka-console-consumer.sh \
--bootstrap-server centos01:9092,centos02:9092,centos03:9092 \
--topic topictest

上述代码中,参数--bootstrap-server用于指定Kafka Broker访问地址。
消费者创建完成后,等待接收生产者的消息。此时在生产者控制台输入消息“hello kafka”后按回车(可以将文件或者标准输入的消息发送到Kafka集群中,默认一行作为一个消息)即可将消息发送到Kafka集群。

file

在消费者控制台,则可以看到输出相同的消息“hello kafka”:

file

谢谢!

Views: 268

初识 Spark (含环境搭建)

大数据开发总体架构

file

什么是Spark

Apache Spark是一个快速通用的集群计算系统,是一种与Hadoop相似的开源集群计算环境,但是Spark在某些工作负载方面表现得更加优越。它提供了Java、Scala、Python和R的高级API,以及一个支持通用的执行图计算的优化引擎。它还支持一组丰富的高级工具,包括使用SQL进行结构化数据处理的Spark SQL、用于机器学习的MLlib、用于图处理的GraphX,以及用于实时流处理的Spark Streaming。

Spark的主要特点:

1、快速
与MapReduce相比,Spark可以支持包括Map和Reduce在内的更多操作,这些操作相互连接形成一个有向无环图(Directed Acyclic Graph,简称DAG),各个操作的中间数据则会被保存在内存中。因此处理速度比MapReduce更加快。Spark通过使用先进的DAG调度器、查询优化器和物理执行引擎,从而能够高性能的实现批处理和流数据处理。

file

2、易用
Spark可以使用Java、Scala、Python、R和SQL快速编写应用程序。
Spark提供了超过80个高级算子(关于算子,在第3章将详细讲解),使用这些算子可以轻松构建并行应用程序,并且可以从Scala、Python、R和SQL的Shell中交互式地使用它们。

3、通用
Spark拥有一系列库,包括SQL和DataFrame、用于机器学习的MLlib、用于图计算的GraphX、用于实时计算的Spark Streaming。可以在同一个应用程序中无缝地组合这些库。

4、到处运行
Spark可以使用独立集群模式运行(使用自带的独立资源调度器,称为Standalone模式),也可以运行在Amazon EC2、Hadoop YARN、Mesos(Apache下的一个开源分布式资源管理框架)、Kubernetes之上,并且可以访问HDFS、Cassandra、HBase、Hive等数百个数据源中的数据。

Spark主要组件

Spark是由多个组件构成的软件栈,Spark 的核心(Spark Core)是一个对由很多计算任务组成的、运行在多个工作机器或者一个计算集群上的应用进行调度、分发以及监控的计算引擎。

file

Spark运行时架构

Spark有多种运行模式,可以运行在一台机器上,称为本地(单机)模式;也可以以YARN或Mesos作为底层资源调度系统以分布式的方式在集群中运行,称为Spark On YARN模式;还可以使用Spark自带的资源调度系统,称为Spark Standalone模式。
本地模式通过多线程模拟分布式计算,通常用于对应用程序的简单测试。本地模式在提交应用程序后,将会在本地生成一个名为“SparkSubmit”的进程,该进程既负责程序的提交又负责任务的分配、执行和监控等。

YARN集群架构

在学习Spark集群架构之前,先需要了解YARN集群的架构。YARN集群总体上是经典的主/从(Master/Slave)架构,主要由ResourceManager、NodeManager、ApplicationMaster和Container等几个组件构成。

file

YARN集群中应用程序的执行流程:

file

Spark Standalone架构

Spark Standalone模式为经典的Master/Slave架构,资源调度是Spark自己实现的。在Standalone模式中,根据应用程序提交的方式不同,Driver(主控进程)在集群中的位置也有所不同。应用程序的提交方式主要有两种:client和cluster,默认是client。

当提交方式为client时,运行架构:

file

当提交方式为cluster时,运行架构:

file

Spark On YARN架构

Spark On YARN模式,遵循YARN的官方规范,YARN只负责资源的管理和调度,运行哪种应用程序由用户自己实现,因此可能在YARN上同时运行MapReduce程序和Spark程序,YARN很好的对每一个程序实现了资源的隔离。这使得Spark与MapReduce可以运行于同一个集群中,共享集群存储资源与计算资源。Spark On YARN模式与Standalone模式一样,也分为client和cluster两种提交方式。

client提交方式架构:

file

cluster提交方式架构:

file

Spark Standalone 集群搭建

Spark Standalone模式的搭建需要在集群的每个节点都安装Spark,集群角色分配如表:

file

1、下载解压安装包
访问Spark官网http://spark.apache.org/downloads.html下载预编译的Spark安装包,选择Spark版本为2.4.0,包类型为“Pre-built for Apache Hadoop 2.7 and later”(Hadoop2.7及之后版本的预编译版本)。

$ tar -zxvf spark-2.4.0-bin-hadoop2.7.tgz -C /opt/modules/

2、修改配置文件
修改slaves文件:

$ cp slaves.template slaves
$ vi slaves

改为以下内容:

centos02
Centos03

修改spark-env.sh文件:

$ cp spark-env.sh.template spark-env.sh
$ vi spark-env.sh

改为以下内容:

export JAVA_HOME=/opt/modules/jdk1.8.0_144
export SPARK_MASTER_HOST=centos01
export SPARK_MASTER_PORT=7077

启动默认的log4j日志配置

cp log4j.properties.template log4j.properties

3、拷贝安装文件到其他节点

$ scp -r /opt/modules/spark-2.4.0-bin-hadoop2.7/ hadoop@centos02:/opt/modules/
$ scp -r /opt/modules/spark-2.4.0-bin-hadoop2.7/ hadoop@centos03:/opt/modules/

4、启动Spark集群
在主节点(centos01)执行,启动Spark集群:

$ sbin/start-all.sh

使用jps查看启动进程,三个节点的进程分别为:Master、Worker、Worker说明启动成功。
访问网址http://centos01:8080,查看Spark的WebUI:

file

Spark提供了一个客户端应用程序提交工具spark-submit,使用该工具可以将编写好的Spark应用程序提交到Spark集群:

$ bin/spark-submit [options] <app jar> [app options]

说明

  • [options]:表示传递给spark-submit的控制参数;
  • <app jar>:表示提交的程序jar包(或Python脚本文件)所在位置;
  • [app options]:表示jar程序需要传递的参数,例如main()方法中需要传递的参数。

例如:

本地模式,使用2个cpu核心:
本地模式不会提交给Spark Master,因此在Spark Master WebUI http://centos01:8080/ 看不到提交的任务信息。

$ bin/spark-submit --master local[2] --deploy-mode client --class org.apache.spark.examples.SparkPi  ./examples/jars/spark-examples_2.11-2.4.8.jar

在Standalone模式(Spark集群使用Spark自带的资源协调服务)下,将Spark自带的求圆周率的程序提交到集群:

$ bin/spark-submit \
--master spark://centos01:7077 \
--class org.apache.spark.examples.SparkPi \
./examples/jars/spark-examples_2.11-2.4.0.jar

说明:
--master参数指定了Master节点的连接地址。该参数根据不同的Spark集群模式,其取值也有所不同:

在输出的日志中间可以找到Pi的估算结果:

Pi is roughly 3.138195690978455

另外在Spark任务运行时,从Spark Master UI界面可以在Running Applications选项下查看当前运行的Spark任务的状态,任务执行完毕后, 在Completed Applications选项中可以找到执行完成的Spark任务,点击Application ID可以查看任务详情和日志信息,

取值 描述
spark://host:port Standalone模式下的Master节点的连接地址,默认端口为7077
yarn 连接到YARN集群。若YARN中没有指定ResourceManager的启动地址,则需要在ResourceManager所在的节点上进行应用程序的提交,否则将因找不到ResourceManager而提交失败
local 运行本地模式,使用1个CPU核心
local[N] 运行本地模式,使用N个CPU核心。例如,local[2]表示使用2个CPU核心运行程序
local[*] 运行本地模式,尽可能使用最多的CPU核心

spark-submit还提供了一些控制资源使用和运行时环境的参数:

参数 描述
--master Master节点的连接地址。取值为spark://host:port、mesos://host:port、yarn、k8s://https://host:port或local (默认为 local[*])
--deploy-mode 提交方式。取值为“client”或“cluster”。“client”表示在本地客户端启动Driver程序,“cluster”表示在集群内部的工作节点上启动Driver程序。默认为“client”
--class 应用程序的主类(Java或Scala程序)
--name 应用程序名称,会在Spark Web UI中显示
--jars 应用依赖的第三方的jar包列表,以逗号分隔
--files 需要放到应用工作目录中的文件列表,以逗号分隔。此参数一般用来放需要分发到各节点的数据文件
--conf 设置任意的SparkConf配置属性。格式为“属性名=属性值”
--properties-file 加载外部包含键值对的属性文件。如果不指定,默认将读取Spark安装目录下的conf/spark-defaults.conf文件中的配置
--driver-memory Driver进程使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1024M
--executor-memory 每个Executor进程所使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1G
--driver-cores 每个Executor进程所使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1G
--executor-cores 每个Executor进程所使用的CPU核心数,默认为1
--num-executors Executor进程数量,默认为2。如果开启动态分配,则初始Executor的数量至少是此参数配置的数量。需要注意的是,此参数仅在Spark On YARN模式中使用

例如,在Standalone模式下,将Spark自带的求圆周率的程序提交到集群,并且设置Driver进程使用内存为512M,每个Executor进程使用内存为1G,每个Executor进程所使用的CPU核心数为1,提交方式为cluster(即Driver进程运行在集群的工作节点中),执行命令:

$ bin/spark-submit \
--master spark://centos01:7077 \
--deploy-mode cluster \
--class org.apache.spark.examples.SparkPi \
--driver-memory 512m \
--executor-memory 1g \
--executor-cores 1 \
./examples/jars/spark-examples_2.11-2.4.0.jar

说明:

查看UI界面可以得知每个Worker的CPU核心数,注意设置的executor-cores不能超过这个数量。

Spark带有交互式的Shell,可在Spark Shell中直接编写Spark任务,然后提交到集群与分布式数据进行交互,并且可以立即查看输出结果。

Spark Standalone模式启动Spark Shell终端:

$ bin/spark-shell --master spark://centos01:7077

Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://hadoop100:4040
Spark context available as 'sc' (master = spark://centos01:7077, app id = app-20210727162737-0001).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.0:
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_281)
Type in expressions to have them evaluated.
Type :help for more information.

scala>

启动完成后,访问Spark Master WebUI http://centos01:8080/ 查看运行的Spark应用程序:

file

退出Spark Shell:(注意命令前面以冒号:开头,可以简写为:q)

scala>:quit

Spark On YARN 集群模式搭建

Spark On YARN模式下Spark Shell的启动与Standalone模式所不同的是:--master的参数值为yarn。例如以下启动命令:

$ bin/spark-shell --master yarn

如果之前没有配置 Spark On YARN 集群模式的环境的话,这一步铁定会遇到异常的,当我们解决这些异常之后, Spark On YARN 集群模式也就自然搭建完成了。

Spark On YARN模式启动Shell出现的问题

Unable to load native-hadoop library for your platform

原因 这个只是警告,而不是错误,提示缺少对Hadoop的lib的引用。在环境变量里面进行设置即可。如果没有配置的话则会使用内建的Java类来实现,导致执行效率上有一定影响。

解决方法

编辑 /etc/profile,添加:

export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native/:$LD_LIBRARY_PATH

使环境变量生效

source /etc/profile

也可以在需要执行的脚本中前面加上

#!/bin/bash
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native/

either HADOOP_CONF_DIR or YARN_CONF_DIR must be set

解决方法
编辑spark-env.sh, 增加

export HADOOP_CONF_DIR=$HADOOP_HOME/etc/hadoop

Neither spark.yarn.jars nor spark.yarn.archive is set

$ hdfs dfs -mkdir -p /user/hadoop/spark/jars
$ hdfs dfs -put /opt/pkg/spark/jars/* /user/hadoop/spark/jars
$ hdfs dfs -chmod -R 755 /user/hadoop/spark/jars

spark-defaults.conf中写入 (注意后面的/*别漏加)

spark.yarn.jars  hdfs://centos01:8020/user/hadoop/spark/jars/*

说明:
centos01:8020 对应HDFS的NameNode的主机名和端口号(对于hadoop2.x的NameNode端口默认是9000,对于hadoop3.x的NameNode端口默认是8020)

最后重新启动Spark,然后运行Spark Shell:

$ bin/spark-shell --master yarn
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://centos01:4040
Spark context available as 'sc' (master = yarn, app id = application_1627371733487_0005).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.8
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_281)
Type in expressions to have them evaluated.
Type :help for more information.

这就是正常运行的样子。

Container Killed on request.

出现这个错误说明Yarn给Spark任务的内存分配过小,Yarn最终直接将Container的进程杀掉了。

解决方法:
在Hadoop的配置文件yarn-site.xml中加入以下内容即可:

    <!--关闭物理内存检查-->
    <property>
        <name>yarn.nodemanager.pmem-check-enabled</name>
        <value>false</value>
    </property>
    <!--关闭虚拟内存检查-->
    <property>
        <name>yarn.nodemanager.vmem-check-enabled</name>
        <value>false</value>
    </property>

修改完毕后,将此文件分发到所有Yarn节点,重启Yarn集群。

Spark 集群 Master 的 HA 环境搭建

默认情况下在Spark standalone集群中进行计算时,由于是RDD的计算模型,所以可以认为worker 已经是有HA特性的了,但是负责资源调度的Master节点有可能出现单点故障。所以为了保证环境的稳定,还是需要配置HA功能。

官方文档中提供了两种HA的机制

  • 基于Zookeeper:利用ZooKeeper来提供主节点选举和集群状态的存储,可以在集群中运行多个连接到同一个Zookeeper集群的Master节点。其中一个将被选为“Leador”和其他的节点将会是Standby模式。如果当前的Leador挂了,会通过选举产生一个新的Leador。
    • 请注意,对于新运行的应用, 延迟可能会有1-2分钟,对于已经运行的应用不会受到影响。
  • 基于文件系统:将状态数据存在目录文件中,当主节点挂掉后,通过重启来解决问题。由于stop-master.sh方式停止Master节点是,不会将对应的数据删除,所以时间长了,可能会影响到启动速度

为了解决Spark的Master单点故障问题,这里基于Zookeeper实现两个Master的主备切换。

首先停止集群,修改spark-env.sh,删除SPARK_MASTER_IP属性配置(如果有),并添加以下配置

export SPARK_DAEMON_JAVA_OPTS="-Dspark.deploy.recoveryMode=ZOOKEEPER -Dspark.deploy.zookeeper.url=centos01:2181,centos02:2181,centos03:2181 -Dspark.deploy.zookeeper.dir=/spark"

然后同步spark-env.sh文件到其他节点

$ scp $SPARK_HOME/conf/spark-env.sh hadoop@centos01:$SPARK_HOME/conf/
spark-env.sh
$ scp $SPARK_HOME/conf/spark-env.sh hadoop@centos02:$SPARK_HOME/conf/
spark-env.sh

在主Master节点启动集群

$ sbin/start-all.sh

在备Master节点启动Master

$ sbin/start-master.sh

异常解决

一些可能会遇到的异常如下:

配置HA时备用Master启动失败

节点centos02启动master后,使用jps命令发现没有master进程,查看logs下面的master日志发现如下报错:
Service 'sparkMaster' failed after 16 retries

解决办法:

在备用Master所在节点的spark-env.sh中配置SPARK_LOCAL_IP属性,对应的值就是此节点的IP或者主机名(Cluster模式下也可以配置localhost或者127.0.0.1, 但是不建议)

# Options read by executors and drivers running inside the cluster
# - SPARK_LOCAL_IP, to set the IP address Spark binds to on this node
export SPARK_LOCAL_IP=centos02

spark master web ui 默认端口 8080 被占用

spark master web ui 默认端口为8080,当系统有其它程序(hadoop3.x版本的集群中的有的节点会启动jetty,用的就是8080端口)也在使用该接口时,启动master时就会报错,为了避免端口冲突,我们也可以自行设置端口号,修改方法:

方法一(推荐)

spark-env.sh中配置SPARK_MASTER_WEBUI_PORT属性,对应的值就是端口号(注意不要和已经开放的端口冲突):

# - SPARK_MASTER_WEBUI_PORT, to use non-default ports for the master
export SPARK_MASTER_WEBUI_PORT=8085
方法二

如果没有在spark-env.sh中指定SPARK Master Web UI端口,则会使用Master启动脚本sbin/start-master.sh中定义的默认端口号,即变量SPARK_MASTER_WEBUI_PORT属性,默认是8080,所以也可以修改这个默认端口号:

if [ "$SPARK_MASTER_WEBUI_PORT" = "" ]; then
  SPARK_MASTER_WEBUI_PORT=8080
fi

现在分别打开活动Master和备用Master的UI界面,可以发现一个状态是Active,另一个状态为Master,说明当前的HA配置成功。

centos01(活动 Master 节点)

URL: spark://centos01:7077
Alive Workers: 2
Cores in use: 2 Total, 0 Used
Memory in use: 4.5 GB Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: ALIVE

centos02(备用 Master 节点)

URL: spark://centos02:7078
Alive Workers: 0
Cores in use: 0 Total, 0 Used
Memory in use: 0.0 B Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: STANDBY

备注:
Master使用的端口是从7077开始,如果找不到就+1, 如此下去直至找到可用端口,默认最多尝试16次,如果还没有找到就报错。

Spark HA测试

接下来进行主备切换测试:

首先centos01关闭Master进程(也可以kill -9直接杀死Master进程)模拟活动Master故障的情况:

此时 http://centos01:8080/ 就不能访问了。

过几秒钟后centos02的Master Web UI界面http://centos02:8080/ 可以发现Master状态变为恢复状态 - RECOVERING,这个状态的持续时间很短,只有几秒钟:

http://centos02:8080/

Spark Master at spark://centos02:7078
URL: spark://centos02:7078
Alive Workers: 0
Cores in use: 0 Total, 0 Used
Memory in use: 0.0 B Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: RECOVERING

然后很快就切换成活动状态 - Active,说明Spark的HA配置成功。
http://centos02:8080/

Spark Master at spark://centos02:7078
URL: spark://centos02:7078
Alive Workers: 2
Cores in use: 2 Total, 0 Used
Memory in use: 4.5 GB Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: ALIVE

Spark 三种提交模式测试

下面的测试环境使用三台虚拟机,系统是CentOS7_x64、主机名分别使hadoop100、hadoop101、Hadoop102,安装了Java JDK8和hadoop3.1.4,Spark 2.4.8 (使用scala2.11编译的) built for Hadoop 2.7.3(用的hadoop的32位的库,因此有些兼容问题),其中Spark的活动master在hadoop100,备用master在hadoop101,worker节点是hadoop101和hadoop102。YARN的ResourceManager节点在hadoop101。

虚拟机环境介绍完了,下面编写一些脚本,使用官方提供的估算Pi值的程序,分别测试Local(本地)、Standalone(提交给Spark集群)以及Spark on YARN(提交给YARN集群)这三种提交模式,除了Local提交模式之外,根据Driver的的部署方式还分为Client模式部署以及Cluster模式部署两种部署方法。

本地客户端提交:
仅使用当前节点计算,默认使用一个CPU核心。
由于客户端和Driver在一个进程,结果直接显示在控制台。

Local[2]表示使用2个核心, Local[*]表示使用尽可能多的核心

#!/bin/bash
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master local \
--name pi-local \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

Standalone-Client部署模式:
提交任务使用spark自带的资源管理机制
计算任务交给Spark配置的工作节点(工作节点配置在slaves中).
Drive和客户端在一个进程,直接在控制台输出结果

#!/bin/bash
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop100:7077 \
--total-executor-cores 2 \
--name pi-cluster-client \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

提交Python编写的Spark任务到Spark集群(虚拟机中的python版本为2.75)
这里使用的是Standalone-Client模式

#!/bin/bash
# Run a Python application on a Spark standalone-client mode
# spark2.4.8预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
  --master spark://hadoop100:7077 \
  /opt/pkg/spark/examples/src/main/python/pi.py \
  100

Standalone-Cluster提交,使用2个CPU核心
Standalone提交任务使用spark集群自带的资源管理机制
cluster部署模式:Master选择一个Workder节点(slaves.sh中配置)运行driver进程
注意这种提交方式要求应用程序使用的jar包和文件需要同步到所有worker节点(或放在HDFS上)
pi的计算结果可以访问Storm WebUI在driver的stdout日志中查找

export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop101:6066 \
--deploy-mode cluster \
--total-executor-cores 2 \
--name pi-cluster-cluster \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

带监督模式的Standalone-Cluster部署

Cluster部署有个好处就是可以开启监督模式(supervise)
开启监督运行模式后当任务失败后可以自动重启

#!/bin/bath
# spark2.4.8预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop100:6066 \
--deploy-mode cluster \
# --supervise 失败时自动重启driver
--supervise \
--total-executor-cores 2 \
--name pi-cluster-cluster \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

spark以client方式提交时,port应该设置为7077;以cluster方式提交时,port设置为6066,因为这种方式提交时,是以rest api方式提交application。

YARN-Client部署模式:
客户端运行driver,任务提交到YARN集群进行调度. 使用YARN集群的工作节点进行计算和(和slaves的配置无关)。可以在YARN的可视化界面(默认端口8088)查看任务执行情况,得到pi的计算结果再发送给Driver在控制台显示

#!/bin/bash
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master yarn \
--deploy-mode client \
--total-executor-cores 2 \
--name pi-yarn-yarn \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

YARN-cluster部署模式
客户端将任务交给YARN集群进行调度
使用YARN集群中的工作节点和spark中的slaves.sh配置无关
driver运行在YARN中的某个NM节点上
可以在YARN的可视化界面(默认端口8088)查看任务执行情况
结果在Application的stdout日志中查看(需提前启动JobHsotry Server)
或者使用yarn logs --applicationId <applicationId>查看

#!/bin/bash
# YARN-cluster部署模式:
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master yarn \
--deploy-mode cluster \
--total-executor-cores 2 \
--name pi-yarn-yarn \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

Spark Shell的使用

Spark Shell的使用也有三种模式

  1. 本地(单机)模式下启动Spark Shell

    即不加--master参数,直接使用bin/spark-shell命令启动 Spark Shell,在本地模式下,所有操作任务只是在本地,也就是当前节点运行,而不会分发到整个集群。

  2. Spark Standalone模式下启动Spark Shell
    使用任意Spark节点进入Spark安装目录,执行以下命令,自动Spark Shell客户端:

$ bin/spark-shell --master spark://hadoop100:7077

说明:
--master指定Master节点的访问地址,centos01为Master所在节点主机名,7077为Master默认端口。

在Spark Shell启动过程日志中可以看出,有个Spark的上下文变量叫做sc,这个变量可以在Spark Shell中直接使用,它也是Spark应用程序的入口,负责于Spark集群进行交互。
启动完成后可以在 http://centos01:8080/ 查看运行的Spark应用程序。

  1. Spark on YARN 模式下启动Spark Shell

前提,需要启动Hadoop的HDFS文件系统和Yarn的相关进程,并完成Spark on YARN 模式的Spark环境搭建。启动命令如下:

$ bin/spark-shell --master yarn

Views: 171

数据同步工具 – DataX (阿里开源)

1、DataX 基本介绍

  • DataX 是阿里巴巴集团内被广泛使用的离线数据同步工具,致力于实现包括:关系型数据库(MySQL、Oracle等)、HDFS、Hive、HBase、ODPS、FTP等各种异构数据源之间稳定高效的数据同步功能。

    datax

  • 设计理念

    • 为了解决异构数据源同步问题,DataX将复杂的网状的同步链路变成了星型数据链路,DataX作为中间传输载体负责连接各种数据源。当需要接入一个新的数据源的时候,只需要将此数据源对接到DataX,便能跟已有的数据源做到无缝数据同步。
  • 当前使用现状

    • DataX在阿里巴巴集团内被广泛使用,承担了所有大数据的离线同步业务,并已持续稳定运行了6年之久。目前每天完成同步8w多道作业,每日传输数据量超过300TB。
    • 此前已经开源DataX1.0版本,此次介绍为阿里云开源全新版本DataX 3.0,有了更多更强大的功能和更好的使用体验。Github主页地址:https://github.com/alibaba/DataX

DataX-logo

2、DataX 3.0 框架设计

  • DataX 本身作为离线数据同步框架,采用Framework + plugin 架构构建。将数据源读取和写入抽象成为Reader/Writer插件,纳入到整个同步框架中。
    • Reader
    • Reader 为数据采集模块,负责采集数据源的数据,将数据发送给 Framework。
    • Writer
    • Writer 为数据写入模块,负责不断向 Framework 取数据,并将数据写入到目的端。
    • Framework
    • Framework 用于连接 Reader 和 Writer,作为两者的数据传输通道,并处理缓冲,流控,并发,数据转换等核心技术问题。

datax3.0框架设计

3、DataX 3.0 插件体系

  • 经过几年积累,DataX目前已经有了比较全面的插件体系,主流的RDBMS数据库、NOSQL、大数据计算系统都已经接入。DataX目前支持数据如下:
类型 数据源 Reader(读) Writer(写) 文档
RDBMS 关系型数据库 MySQL
Oracle
SQLServer
PostgreSQL
DRDS
达梦
通用RDBMS(支持所有关系型数据库)
阿里云数仓数据存储 ODPS
ADS
OSS
OCS
NoSQL数据存储 OTS
Hbase0.94
Hbase1.1
MongoDB
Hive
无结构化数据存储 TxtFile
FTP
HDFS
Elasticsearch

4、DataX 3.0 核心架构

  • DataX 3.0 支持单机多线程模式完成 数据同步作业,本小节按一个DataX作业生命周期的时序图,从整体架构设计简要说明DataX各个模块之间的相互关系。

    DataX3.0核心架构

  • 核心模块介绍

      1. DataX完成单个数据同步的作业,我们称之为Job,DataX接受到一个Job之后,将启动一个进程来完成整个作业同步过程。DataX Job模块是单个作业的中枢管理节点,承担了数据清理、子任务切分(将单一作业计算转化为多个子Task)、TaskGroup管理等功能。
      1. DataX Job启动后,会根据不同的源端切分策略,将Job切分成多个小的Task(子任务),以便于并发执行。Task便是DataX作业的最小单元,每一个Task都会负责一部分数据的同步工作。
      1. 切分多个Task之后,DataX Job会调用 Scheduler 模块,根据配置的并发数据量,将拆分成的Task重新组合,组装成TaskGroup(任务组)。每一个TaskGroup负责以一定的并发运行完毕分配好的所有Task,默认单个任务组的并发数量为5。
      1. 每一个Task都由TaskGroup负责启动,Task启动后,会固定启动Reader—>Channel—>Writer的线程来完成任务同步工作。
      1. DataX作业运行起来之后, Job监控并等待多个TaskGroup模块任务完成,等待所有TaskGroup任务完成后Job成功退出。否则,异常退出,进程退出值非0。
  • DataX调度流程

  • 举例来说,用户提交了一个DataX作业,并且配置了20个并发,目的是将一个100张分表的mysql数据同步到odps里面。 DataX的调度决策思路是:

      1. DataXJob根据分库分表切分成了100个Task。
      1. 根据20个并发,默认单个任务组的并发数量为5,DataX计算共需要分配4个TaskGroup。
      1. 这里4个TaskGroup平分切分好的100个Task,每一个TaskGroup负责以5个并发共计运行25个Task。

5、DataX 安装部署

  • 安装前置要求

    • Linux
    • JDK ( 1.8 以上 )
    • Python ( 2.6 以上 )
  • 1、访问官网下载安装包

  • 2、上传安装包到服务器node01节点

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

    tar -zxvf datax.tar.gz -C /kkb/install
  • 4、运行自检脚本测试

    [hadoop@node01 bin]$ cd /kkb/install/datax/bin
    [hadoop@node01 bin]$ python datax.py ../job/job.json 

    image-20210511162206902

6、DataX 实战案例

6.1 从stream流读取数据并打印到控制台

  • 需求:使用datax实现读取字符串,然后打印到控制台。

  • 1、创建作业的配置文件(json格式)

    • 可以通过命令查看配置模板: python datax.py -r {YOUR_READER} -w {YOUR_WRITER}

    • 查看配置模板,执行脚本命令

    [hadoop@node01 datax]$ cd /kkb/install/datax
    [hadoop@node01 datax]$ python bin/datax.py -r streamreader -w streamwriter
    ##查看输出结果
    
    DataX (DATAX-OPENSOURCE-3.0), From Alibaba !
    Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
    
    Please refer to the streamreader document:
         https://github.com/alibaba/DataX/blob/master/streamreader/doc/streamreader.md 
    
    Please refer to the streamwriter document:
         https://github.com/alibaba/DataX/blob/master/streamwriter/doc/streamwriter.md 
    
    Please save the following configuration as a json file and  use
         python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json 
    to run the job.
    
    {
        "job": {
            "content": [
                {
                    "reader": {
                        "name": "streamreader", 
                        "parameter": {
                            "column": [], 
                            "sliceRecordCount": ""
                        }
                    }, 
                    "writer": {
                        "name": "streamwriter", 
                        "parameter": {
                            "encoding": "", 
                            "print": true
                        }
                    }
                }
            ], 
            "setting": {
                "speed": {
                    "channel": ""
                }
            }
        }
    }
  • 2、根据模板写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 stream2stream.json, 文件内容如下:
    {
      "job": {
        "content": [
          {
            "reader": {
              "name": "streamreader",
              "parameter": {
                "sliceRecordCount": 10,
                "column": [
                  {
                    "type": "long",
                    "value": "10"
                  },
                  {
                    "type": "string",
                    "value": "hello,你好,世界-DataX"
                  }
                ]
              }
            },
            "writer": {
              "name": "streamwriter",
              "parameter": {
                "encoding": "UTF-8",
                "print": true
              }
            }
          }
        ],
        "setting": {
          "speed": {
            "channel": 5,
            "bytes":0
           },
           "errorLimit": {
             "record": 10,
             "percentage": 0.02
            }
        }
      }
    }
    • 其中sliceRecordCount表示每个channel 生成数据的条数。

    • speed表示限速

    • channel表示任务并发数。

    • bytes表示每秒字节数,默认为0(不限速)。

    • errorLimit表示错误控制

    • record: 出错记录数超过record设置的条数时,任务标记为失败.

    • percentage: 当出错记录数超过percentage百分数时,任务标记为失败.

  • 3、动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/stream2stream.json 
  • 4、观察控制台输出结果

    同步结束,显示日志如下:
    
    10      hello,你好,世界-DataX
    10      hello,你好,世界-DataX
    10      hello,你好,世界-DataX
    10      hello,你好,世界-DataX
    10      hello,你好,世界-DataX
    10      hello,你好,世界-DataX
    
    ...
    2021-05-11 16:52:39.274 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-11 16:52:29
    任务结束时刻                    : 2021-05-11 16:52:39
    任务总计耗时                    :                 10s
    任务平均流量                    :               95B/s
    记录写入速度                    :              5rec/s
    读出记录总数                    :                  50
    读写失败总数                    :                   0

6.2 从mysql表读取数据并打印到控制台

  • 需求:使用datax实现读取mysql一张表指定字段的数据,打印到控制台

  • 1、在mysql数据库中创建student表,并且加载数据到表中

    mysql> create database datax;
    mysql> use datax;
    mysql> create table student(id int,name varchar(20),age int,createtime timestamp );
    mysql> insert into <code>student (id, name, age, createtime) values('1','zhangsan','18','2021-05-10 18:10:00');
    
    mysql> insert into student (id, name, age, createtime) values('2','lisi','28','2021-05-10 19:10:00');
    
    mysql> insert into student (id, name, age, createtime) values('3','wangwu','38','2021-05-10 20:10:00');
  • 2、创建作业的配置文件(json格式)

    • 执行脚本命令(查看配置模板)
    [hadoop@node01 datax]$ cd /kkb/install/datax
    [hadoop@node01 datax]$ python bin/datax.py -r mysqlreader -w streamwriter
    
    DataX (DATAX-OPENSOURCE-3.0), From Alibaba !
    Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
    
    Please refer to the mysqlreader document:
         https://github.com/alibaba/DataX/blob/master/mysqlreader/doc/mysqlreader.md 
    
    Please refer to the streamwriter document:
         https://github.com/alibaba/DataX/blob/master/streamwriter/doc/streamwriter.md 
    
    Please save the following configuration as a json file and  use
         python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json 
    to run the job.
    
    {
        "job": {
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader", 
                        "parameter": {
                            "column": [], 
                            "connection": [
                                {
                                    "jdbcUrl": [], 
                                    "table": []
                                }
                            ], 
                            "password": "", 
                            "username": "", 
                            "where": ""
                        }
                    }, 
                    "writer": {
                        "name": "streamwriter", 
                        "parameter": {
                            "encoding": "", 
                            "print": true
                        }
                    }
                }
            ], 
            "setting": {
                "speed": {
                    "channel": ""
                }
            }
        }
    }
  • 3、根据模板写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 mysql2stream.json, 文件内容如下:
    {
        "job": {
            "setting": {
                "speed": {
                     "channel": 3
                },
                "errorLimit": {
                    "record": 0,
                    "percentage": 0.02
                }
            },
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader",
                        "parameter": {
                            "username": "root",
                            "password": "123456",
                            "column": [
                                "id",
                                "name",
                              "age",
                                "createtime"
                            ],
                            "connection": [
                                {
                                    "table": [
                                        "student"
                                    ],
                                    "jdbcUrl": [
         "jdbc:mysql://node03:3306/datax"
                                    ]
                                }
                            ]
                        }
                    },
                   "writer": {
                        "name": "streamwriter",
                        "parameter": {
                            "print":true
                        }
                    }
                }
            ]
        }
    }
    
  • 4、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/mysql2stream.json 
  • 5、观察控制台输出结果

    同步结束,显示日志如下:
    
    1       zhangsan        18      2021-05-10 18:10:00
    2       lisi    28      2021-05-10 19:10:00
    3       wangwu  38      2021-05-10 20:10:00
    
    ...
    2021-05-11 17:31:29.904 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-11 17:31:19
    任务结束时刻                    : 2021-05-11 17:31:29
    任务总计耗时                    :                 10s
    任务平均流量                    :                2B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   3
    读写失败总数                    :                   0

6.3 从mysql表读取增量数据并打印到控制台

  • 需求:使用datax实现mysql表增量数据同步打印到控制台。

  • 1、创建作业的配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 mysqlAdd2stream.json, 文件内容如下:
    {
      "job": {
          "setting": {
              "speed": {
                   "channel": 3
              },
              "errorLimit": {
                  "record": 10,
                  "percentage": 0.02
              }
          },
          "content": [
              {
                  "reader": {
                      "name": "mysqlreader",
                      "parameter": {
                          "username": "root",
                          "password": "123456",
                          "column": [
                              "id",
                              "name",
                          "age",
                              "createtime"
                          ],
                        "where":"createtime > '${start_time}' and createtime < '${end_time}'",
                          "connection": [
                              {
                                  "table": [
                                      "student"
                                  ],
                                  "jdbcUrl": [
       "jdbc:mysql://node03:3306/datax"
                                  ]
                              }
                          ]
                      }
                  },
                 "writer": {
                      "name": "streamwriter",
                      "parameter": {
                          "print":true
                      }
                  }
              }
          ]
      }
    }
    
  • 2、向student表中插入一条数据

    mysql> insert into <code>student (id, name, age, createtime) values('4','xiaoming','48','2021-05-11 19:10:00')
  • 3、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/mysqlAdd2stream.json -p "-Dstart_time='2021-05-11 00:00:00' -Dend_time='2021-05-11 23:59:59'" 
  • 4、观察控制台输出结果

    同步结束,显示日志如下:
    ...
    INFO  CommonRdbmsReader$Task - Finished read record by Sql: [select id,name,age,createtime from student where (createtime > '2021-05-11 00:00:00' and createtime < '2021-05-11 23:59:59')
    
    4       xiaoming        48      2021-05-11 19:10:00
    
    ...
    2021-05-11 18:37:35.755 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-11 18:37:25
    任务结束时刻                    : 2021-05-11 18:37:35
    任务总计耗时                    :                 10s
    任务平均流量                    :                1B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   1
    读写失败总数                    :                   0

6.4 使用datax实现mysql2mysql

  • 需求:使用datax实现将数据从mysql当中读取,并且通过sql语句实现数据的过滤,并且将数据写入到mysql另外一张表当中去。

  • 1、创建作业的配置文件(json格式)

    • 查看配置模板,执行脚本命令
    [hadoop@node01 datax]$ cd /kkb/install/datax
    [hadoop@node01 datax]$ python bin/datax.py -r mysqlreader -w mysqlwriter
    
    DataX (DATAX-OPENSOURCE-3.0), From Alibaba !
    Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
    
    Please refer to the mysqlreader document:
         https://github.com/alibaba/DataX/blob/master/mysqlreader/doc/mysqlreader.md 
    
    Please refer to the mysqlwriter document:
         https://github.com/alibaba/DataX/blob/master/mysqlwriter/doc/mysqlwriter.md 
    
    Please save the following configuration as a json file and  use
         python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json 
    to run the job.
    
    {
        "job": {
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader", 
                        "parameter": {
                            "column": [], 
                            "connection": [
                                {
                                    "jdbcUrl": [], 
                                    "table": []
                                }
                            ], 
                            "password": "", 
                            "username": "", 
                            "where": ""
                        }
                    }, 
                    "writer": {
                        "name": "mysqlwriter", 
                        "parameter": {
                            "column": [], 
                            "connection": [
                                {
                                    "jdbcUrl": "", 
                                    "table": []
                                }
                            ], 
                            "password": "", 
                            "preSql": [], 
                            "session": [], 
                            "username": "", 
                            "writeMode": ""
                        }
                    }
                }
            ], 
            "setting": {
                "speed": {
                    "channel": ""
                }
            }
        }
    }
  • 2、根据模板写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 mysql2mysql.json, 文件内容如下:
    {
        "job": {
            "setting": {
                "speed": {
                     "channel":1
                }
            },
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader",
                        "parameter": {
                            "username": "root",
                            "password": "123456",
                            "connection": [
                                {
                                    "querySql": [
                                        "select id,name,age,createtime from student where age < 30;"
                                    ],
                                    "jdbcUrl": [
                                        "jdbc:mysql://node03:3306/datax"
                                    ]
                                }
                            ]
                        }
                    },
                      "writer": {
                        "name": "mysqlwriter",
                        "parameter": {
                            "writeMode": "insert",
                            "username": "root",
                            "password": "123456",
                            "column": [
                                "id",
                                "name",
                                "age",
                                "createtime"
                            ],
                            "preSql": [
                                "delete from person"
                            ],
                            "connection": [
                                {
                                    "jdbcUrl": "jdbc:mysql://node03:3306/datax?useUnicode=true&characterEncoding=utf-8",
                                    "table": [
                                        "person"
                                    ]
                                }
                            ]
                        }
                    }
                }
            ]
        }
    }
    
  • 3、创建目标表

    mysql> create table datax.person(id int,name varchar(20),age int,createtime timestamp );
  • 4、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/mysql2mysql.json 
  • 5、观察控制台输出结果

    同步结束,显示日志如下:
    
    2021-05-12 11:17:24.390 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-12 11:17:13
    任务结束时刻                    : 2021-05-12 11:17:24
    任务总计耗时                    :                 10s
    任务平均流量                    :                3B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   2
    读写失败总数                    :                   0
  • 6、查看person表数据

    image-20210512111811791

6.5 使用datax实现将mysql数据导入到hdfs

  • 需求: 将mysql表student的数据导入到hdfs的 /datax/mysql2hdfs/ 路径下面去。

  • 1、创建作业的配置文件(json格式)

    • 执行脚本命令查看配置模板
    [hadoop@node01 datax]$ cd /kkb/install/datax
    [hadoop@node01 datax]$ python bin/datax.py -r mysqlreader -w hdfswriter
    
    DataX (DATAX-OPENSOURCE-3.0), From Alibaba !
    Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
    
    Please refer to the mysqlreader document:
         https://github.com/alibaba/DataX/blob/master/mysqlreader/doc/mysqlreader.md 
    
    Please refer to the hdfswriter document:
         https://github.com/alibaba/DataX/blob/master/hdfswriter/doc/hdfswriter.md 
    
    Please save the following configuration as a json file and  use
         python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json 
    to run the job.
    
    {
        "job": {
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader", 
                        "parameter": {
                            "column": [], 
                            "connection": [
                                {
                                    "jdbcUrl": [], 
                                    "table": []
                                }
                            ], 
                            "password": "", 
                            "username": "", 
                            "where": ""
                        }
                    }, 
                    "writer": {
                        "name": "hdfswriter", 
                        "parameter": {
                            "column": [], 
                            "compress": "", 
                            "defaultFS": "", 
                            "fieldDelimiter": "", 
                            "fileName": "", 
                            "fileType": "", 
                            "path": "", 
                            "writeMode": ""
                        }
                    }
                }
            ], 
            "setting": {
                "speed": {
                    "channel": ""
                }
            }
        }
    }
  • 2、根据模板写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 mysql2hdfs.json, 文件内容如下:
    {
        "job": {
            "setting": {
                "speed": {
                     "channel":1
                }
            },
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader",
                        "parameter": {
                            "username": "root",
                            "password": "123456",
                            "connection": [
                                {
                                    "querySql": [
                                        "select id,name,age,createtime from student where age < 30;"
                                    ],
                                    "jdbcUrl": [
                                        "jdbc:mysql://node03:3306/datax"
                                    ]
                                }
                            ]
                        }
                    },
                      "writer": {
                        "name": "hdfswriter",
                        "parameter": {
                            "defaultFS": "hdfs://node01:8020",
                            "fileType": "text",
                            "path": "/datax/mysql2hdfs/",
                            "fileName": "student.txt",
                            "column": [
                                {
                                    "name": "id",
                                    "type": "INT"
                                },
                                {
                                    "name": "name",
                                    "type": "STRING"
                                },
                                {
                                    "name": "age",
                                    "type": "INT"
                                },
                                {
                                    "name": "createtime",
                                    "type": "TIMESTAMP"
                                }
                            ],
                            "writeMode": "append",
                            "fieldDelimiter": "\t",
                            "compress":"gzip"
                        }
                    }
                }
            ]
        }
    }
    
  • 3、启HDFS, 创建目标路径

    [hadoop@node01 ~]$ start-dfs.sh 
    [hadoop@node01 ~]$ hdfs dfs -mkdir -p /datax/mysql2hdfs
  • 4、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/mysql2hdfs.json 
  • 5、观察控制台输出结果

    同步结束,显示日志如下:
    
    2021-05-12 11:32:26.452 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-12 11:32:14
    任务结束时刻                    : 2021-05-12 11:32:26
    任务总计耗时                    :                 11s
    任务平均流量                    :                3B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   2
    读写失败总数                    :                   0
  • 6、查看HDFS上文件生成

    image-20210512113442247

6.6 使用datax实现将hdfs数据导入到mysql表中

  • 需求: 将hdfs上数据文件 user.txt 导入到mysql数据库的user表中。

  • 1、创建作业的配置文件(json格式)

    • 查看配置模板,执行脚本命令
    [hadoop@node01 datax]$ cd /kkb/install/datax
    [hadoop@node01 datax]$ python bin/datax.py -r hdfsreader -w mysqlwriter
    
    DataX (DATAX-OPENSOURCE-3.0), From Alibaba !
    Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
    
    Please refer to the hdfsreader document:
         https://github.com/alibaba/DataX/blob/master/hdfsreader/doc/hdfsreader.md 
    
    Please refer to the mysqlwriter document:
         https://github.com/alibaba/DataX/blob/master/mysqlwriter/doc/mysqlwriter.md 
    
    Please save the following configuration as a json file and  use
         python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json 
    to run the job.
    
    {
        "job": {
            "content": [
                {
                    "reader": {
                        "name": "hdfsreader", 
                        "parameter": {
                            "column": [], 
                            "defaultFS": "", 
                            "encoding": "UTF-8", 
                            "fieldDelimiter": ",", 
                            "fileType": "orc", 
                            "path": ""
                        }
                    }, 
                    "writer": {
                        "name": "mysqlwriter", 
                        "parameter": {
                            "column": [], 
                            "connection": [
                                {
                                    "jdbcUrl": "", 
                                    "table": []
                                }
                            ], 
                            "password": "", 
                            "preSql": [], 
                            "session": [], 
                            "username": "", 
                            "writeMode": ""
                        }
                    }
                }
            ], 
            "setting": {
                "speed": {
                    "channel": ""
                }
            }
        }
    }
  • 2、根据模板写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 hdfs2mysql.json, 文件内容如下:
    {
        "job": {
            "setting": {
                "speed": {
                     "channel":1
                }
            },
            "content": [
                {
                    "reader": {
                        "name": "hdfsreader",
                        "parameter": {
                          "defaultFS": "hdfs://node01:8020",
                            "path": "/user.txt",                  
                            "fileType": "text",
                            "encoding": "UTF-8",
                            "fieldDelimiter": "\t",
                            "column": [
                                   {
                                    "index": 0,
                                    "type": "long"
                                   },
                                   {
                                    "index": 1,
                                    "type": "string"
                                   },
                                   {
                                    "index": 2,
                                    "type": "long"
                                   }
                            ]
                          }
                      },
                   "writer": {
                        "name": "mysqlwriter",
                        "parameter": {
                            "writeMode": "insert",
                            "username": "root",
                            "password": "123456",
                            "column": [
                                "id",
                                "name",
                              "age"
                            ],
                            "preSql": [
                                "delete from user"
                            ],
                            "connection": [
                                {
                                    "jdbcUrl": "jdbc:mysql://node03:3306/datax?useUnicode=true&characterEncoding=utf-8",
                                    "table": [
                                        "user"
                                    ]
                                }
                            ]
                        }
                    }
                }
            ]
        }
    }
    
  • 3、准备HDFS上测试数据文件 user.txt

    • user.txt文件内容如下
    1   zhangsan    20
    2   lisi    29
    3   wangwu  25
    4   zhaoliu 35
    5   kobe    40
    • 文件中每列字段通过\t 制表符进行分割,上传文件到hdfs上
    [hadoop@node01 ~]$ hdfs dfs -put user.txt /
  • 4、创建目标表

    mysql> create table datax.user(id int,name varchar(20),age int);
  • 5、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/hdfs2mysql.json 
  • 6、观察控制台输出结果

    同步结束,显示日志如下:
    
    任务启动时刻                    : 2021-05-12 12:02:47
    任务结束时刻                    : 2021-05-12 12:02:58
    任务总计耗时                    :                 11s
    任务平均流量                    :                4B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   5
    读写失败总数                    :                   0
  • 7、查看user表数据

    image-20210512120344118

6.7 使用datax实现将mysql数据同步到hive表中

  • 需求 :使用datax将mysql中的 user表数据全部同步到hive表中

  • 1、创建一张hive表

    • 启动hiveserver2
    [hadoop@node03 hive]$ hiveserver2   
    • 通过beeline连接hiveserver2
    [hadoop@node03 hive]$ beeline   
    
    beeline> !connect jdbc:hive2://node03:10000
    • 创建数据库和表
    0: jdbc:hive2://node03:10000> create database datax;
    0: jdbc:hive2://node03:10000> use datax;
    0: jdbc:hive2://node03:10000> create external table t_user(id int,name string,age int) row format delimited fields terminated by '\t';
  • 2、编写配置文件

    • 进入到 /kkb/install/datax/job 目录,然后创建配置文件 mysql2hive.json, 文件内容如下:
    {
        "job": {
            "setting": {
                "speed": {
                     "channel":1
                }
            },
            "content": [
                {
                    "reader": {
                        "name": "mysqlreader",
                        "parameter": {
                            "username": "root",
                            "password": "123456",
                            "connection": [
                                {
                                    "jdbcUrl": [
                                        "jdbc:mysql://node03:3306/datax"
                                    ],
                                    "table": [
                                        "user"
                                    ]
                                }
                            ],
                           "column": [
                                "id",
                                "name",
                              "age"
                            ]
                        }
                    },
                      "writer": {
                        "name": "hdfswriter",
                        "parameter": {
                            "defaultFS": "hdfs://node01:8020",
                            "fileType": "text",
                            "path": "/user/hive/warehouse/datax.db/t_user",
                            "fileName": "user.txt",
                            "column": [
                                {
                                    "name": "id",
                                    "type": "INT"
                                },
                                {
                                    "name": "name",
                                    "type": "STRING"
                                },
                                {
                                    "name": "age",
                                    "type": "INT"
                                }
                            ],
                            "writeMode": "append",
                            "fieldDelimiter": "\t",
                            "compress":"gzip"
                        }
                    }
                }
            ]
        }
    }
    
  • 3、启动DataX

    [hadoop@node01 bin]$ cd /kkb/install/datax
    [hadoop@node01 bin]$ python bin/datax.py job/mysql2hive.json 
  • 4、观察控制台输出结果

    同步结束,显示日志如下:
    
    2021-05-12 12:20:31.080 [job-0] INFO  JobContainer - 
    任务启动时刻                    : 2021-05-12 12:20:19
    任务结束时刻                    : 2021-05-12 12:20:31
    任务总计耗时                    :                 11s
    任务平均流量                    :                4B/s
    记录写入速度                    :              0rec/s
    读出记录总数                    :                   5
    读写失败总数                    :                   0
  • 5、查看hive中t_user表数据

    image-20210512172948034

Views: 116

分布式存储引擎 – 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

Index