Spark SQL 结构化数据处理引擎

什么是Spark SQL

Spark SQL是一个用于结构化数据处理的Spark组件。所谓结构化数据,是指具有Schema信息的数据,例如json、parquet、avro、csv格式的数据。与基础的Spark RDD API不同,Spark SQL提供了对结构化数据的查询和计算接口。

Spark SQL的主要特点:

  1. 将SQL查询与Spark应用程序无缝组合
    Spark SQL允许使用SQL在Spark程序中查询结构化数据。与Hive不同的是,Hive是将SQL翻译成MapReduce作业,底层是基于MapReduce;而Spark SQL底层使用的是Spark RDD。例如以下代码,在Spark应用程序中嵌入SQL语句:

    results = spark.sql( "SELECT * FROM people")
  2. 以相同的方式连接到多种数据源
    Spark SQL提供了访问各种数据源的通用方法,数据源包括Hive、Avro、Parquet、ORC、JSON、JDBC等。例如以下代码,读取HDFS中的JSON文件,然后将该文件的内容创建为临时视图,最后与其他表根据指定的字段关联查询:

    //读取JSON文件
    val userScoreDF = spark.read.json("hdfs://centos01:9000/people.json")
    //创建临时视图user_score
    userScoreDF.createTempView("user_score")
    //根据name关联查询
    val resDF=spark.sql("SELECT i.age,i.name,c.score FROM user_info i " +
                    "JOIN user_score c ON i.name=c.name")
  3. 在现有的数据仓库上运行SQL或HiveQL查询
    Spark SQL支持HiveQL语法以及Hive SerDes和UDF(用户自定义函数),允许访问现有的Hive仓库。

DataFrame和Dataset

DataFrame是Spark SQL提供的一个编程抽象,与RDD类似,也是一个分布式的数据集合。但与RDD不同的是,DataFrame的数据都被组织到有名字的列中,就像关系型数据库中的表一样。此外,多种数据都可以转化为DataFrame,例如:Spark计算过程中生成的RDD、结构化数据文件、Hive中的表、外部数据库等。
DataFrame在RDD的基础上添加了数据描述信息(Schema,即元信息),因此看起来更像是一张数据库表。

file

Spark SQL基本使用

Spark Shell启动时除了默认创建一个名为“sc”的SparkContext的实例外,还创建了一个名为“spark”的SparkSession实例,该“spark”变量也可以在Spark Shell中直接使用。
SparkSession只是在SparkContext基础上的封装,应用程序的入口仍然是SparkContext。SparkSession允许用户通过它调用DataFrame和Dataset相关API来编写Spark程序,支持从不同的数据源加载数据,并把数据转换成DataFrame,然后使用SQL语句来操作DataFrame数据。
例如,在HDFS中有一个文件/input/person.txt,文件内容:

1,zhangsan,25
2,lisi,22
3,wangwu,30

现需要使用Spark SQL将该文件中的数据按照年龄降序排列,步骤:

  1. 加载数据为Dataset
    调用SparkSession的API read.textFile()可以读取指定路径中的文件内容,并加载为一个Dataset:

    $ spark-shell --master yarn
    scala> val d1=spark.read.textFile("hdfs://centos01:9000/input/person.txt")
    d1: org.apache.spark.sql.Dataset[String] = [value: string]

    从变量d1的类型可以看出,textFile()方法将读取的数据转为了Dataset。除了textFile()方法读取文本内容外,还可以使用csv()、jdbc()、json()等方法读取csv文件、jdbc数据源、json文件等数据。调用Dataset中的show()方法可以输出Dataset中的数据内容。查看d1中的数据内容:

    scala> d1.show()
    +-------------+
    |        value|
    +-------------+
    |1,zhangsan,25|
    |2,lisi,22|
    |3,wangwu,30|
    +-------------+

    从上述内容可以看出,Dataset将文件中的每一行看做一个元素,并且所有元素组成了一列,列名默认为“value”。

    如果Spark报错ClassNotFoundException: Class com.hadoop.compression.lzo.Lzo说明Hadoop支持Lzo压缩但是Spark没有配置支持压缩相关的类库,可以修改spark-defaults.conf并添加两行:
    spark.driver.extraClassPath /opt/pkg/hadoop/share/hadoop/common/hadoop-lzo-0.4.20.jar spark.executor.extraClassPath /opt/pkg/hadoop/share/hadoop/common/hadoop-lzo-0.4.20.jar

  2. 给Dataset添加元数据信息
    定义一个样例类Person,用于存放数据描述信息(Schema):

    scala> case class Person(id:Int,name:String,age:Int)

    导入SparkSession的隐式转换,以便后续可以使用Dataset的算子:

    scala> import spark.implicits._

    调用Dataset的map()算子将每一个元素拆分并存入Person类中:

    scala> val personDataset=d1.map(line=>{
         | val fields = line.split(",")
         | val id = fields(0).toInt
         | val name = fields(1)
         | val age = fields(2).toInt
         | Person(id, name, age)
         | })

    此时查看personDataset中的数据内容,personDataset中的数据类似于一张关系型数据库的表:

    scala> personDataset.show()
    +---+--------+---+
    | id|    name|age|
    +---+--------+---+
    |  1|zhangsan| 25|
    |  2|    lisi| 22|
    |  3|  wangwu| 30|
    +---+--------+---+
  3. 将Dataset转为DataFrame
    Spark SQL查询的是DataFrame中的数据,因此需要将存有元数据信息的Dataset转为DataFrame。
    调用Dataset的toDF()方法,将存有元数据的Dataset转为DataFrame,代码:

    scala> val pdf = personDataset.toDF()
    pdf: org.apache.spark.sql.DataFrame = [id: int, name: string ... 1 more field]
  4. 执行SQL查询
    在DataFrame上创建一个临时视图“v_person”,代码:

    scala> pdf.createTempView("v_person")

    使用SparkSession对象执行SQL查询,代码:

    scala> val result = spark.sql("select * from v_person order by age desc")
    result: org.apache.spark.sql.DataFrame = [id: int, name: string ... 1 more field]

    调用show()方法输出结果数据,代码:

    scala> result.show()
    +---+--------+---+
    | id|    name|age|
    +---+--------+---+
    |  3|  wangwu| 30|
    |  1|zhangsan| 25|
    |  2|    lisi | 22|
    +---+--------+---+

    可以看到,结果数据已按照age字段降序排列。

  5. 指定导出文件格式

    # 导出json格式文件
    scala> result.write.format("json").save("/output/json")
    # 导出scv格式文件
    scala> result.write.format("csv").save("/output/csv")
    # 导出orc格式文件
    scala> result.write.format("orc").save("/output/orc")
    # 导出parquet格式文件
    scala> result.write.format("parquet").save("/output/parquet")

    查看导出的文件(有几个文件说明有几个分区):

    $ hadoop fs -cat /output/json/part-*
    {"id":3,"name":"wangwu","age":30}
    {"id":1,"name":"zhangsan","age":25}
    {"id":2,"name":"lisi","age":22}
    
    $ hadoop fs -cat /output/csv/part-*
    3,wangwu,30
    1,zhangsan,25
    2,lisi,22
    
    $ hadoop fs -cat /output/orc/part-*
    ORC...(二进制乱码)
    
    [hadoop@hadoop100 conf]$ hadoop fs -cat /output/parquet/part-*
    PAR1...(二进制乱码)
  6. 分区设置
    Spark中使用SparkSql进行shuffle操作,默认分区数是200个;参数配置是--conf spark.sql.shuffle.partitions, 如果想要修改SparkSQL执行shuffle操作时分区数:

    1. 配置 spark.sql.shuffle.partitions,适用场景spark.sql()合并分区

      spark.conf.set("spark.sql.shuffle.partitions", 5) #后面的数字是你希望的分区数

      这样配置后,通过spark.sql()执行后写出的数据分区数就是你要求的个数,如这里5。

    2. 配置 coalesce(n),适用场景spark写出数据到指定路径下合并分区,不会引起shuffle

      df = spark.sql(sql_string).coalesce(1) #合并分区数
      df.write.format("csv")
      .mode("overwrite")
      .option("sep", ",")
      .option("header", True)
      .save(hdfs_path)
    3. 配置repartition(n), 重新分区,会引发shuffle

      df = spark.sql(sql_string).repartition(1) #重新分区,会引发全局shuffle
      df.write.format("csv")
      .mode("overwrite")
      .option("sep", ",")
      .option("header", True)
      .save(hdfs_path)
  7. 分区分桶导出

    1. 并行写出之 partitionBy() 指定分区列 , 会根据分区列创建子文件夹,并行写出数据
      df.write.mode("overwrite")
      .partitionBy("day")
      .save("/tmp/partitioned-files.parquet")
    2. 并行写出之 repartition() ,一般spark中有几个分区就会有几个并行的IO写出
      df.repartition(5)
      .write.format("csv")
      .save("/tmp/multiple.csv")
    3. 分桶写出,好处是后续读入的时候数据就不会做shuffle了,因为相同分桶的数据会被划分到同一个物理分区中
      csvFile.write.format("parquet")
      .mode("overwrite")
      .bucketBy(5, "gmv") #第一个参数:分成几个桶,第二个参数:按哪列进行分桶
      .saveAsTable("bucketedFiles")

      Spark IO相关API请参考:Spark官方API文档Input and Output

Spark SQL数据源

Spark SQL支持通过DataFrame接口对各种数据源进行操作。DataFrame可以使用相关转换算子进行操作,也可以用于创建临时视图。将DataFrame注册为临时视图可以对其中的数据使用SQL查询。
Spark SQL提供了两个常用的加载数据和写入数据的方法:load()方法和save()方法。load()方法可以加载外部数据源为一个DataFrame,save()方法可以将一个DataFrame写入到指定的数据源。

1、默认数据源
默认情况下,load()方法和save()方法只支持Parquet格式的文件,也可以在配置文件中通过参数spark.sql.sources.default对默认文件格式进行更改。
Spark SQL可以很容易的读取Parquet文件并将其数据转为DataFrame数据集。例如,读取HDFS中的文件/users.parquet,并将其中的name列与favorite_color列写入HDFS的/result目录,代码:

val spark = SparkSession.builder() //创建或得到SparkSession
  .appName("SparkSQLDataSource")
  .master("local[*]")
  .getOrCreate()
//加载parquet格式的文件,返回一个DataFrame集合
val usersDF = spark.read.load("hdfs://centos01:9000/users.parquet")
usersDF.show()
// +------+--------------+----------------+
// |  name|favorite_color|favorite_numbers|
// +------+--------------+----------------+
// |Alyssa|            null|  [3, 9, 15, 20]|
// |   Ben|              red|                []|
// +------+--------------+----------------+
 //查询DataFrame中的name列和favorite_color列,并写入HDFS
usersDF.select("name","favorite_color")
  .write.save("hdfs://centos01:9000/result")

除了使用select()方法查询外,也可以使用SparkSession对象的sql()方法执行SQL语句进行查询,该方法的返回结果仍然是一个DataFrame。

//创建临时视图
usersDF.createTempView("t_user")
//执行SQL查询,并将结果写入到HDFS
spark.sql("SELECT name,favorite_color FROM t_user")
  .write.save("hdfs://centos01:9000/result")

2、手动指定数据源

使用format()方法可以手动指定数据源。数据源需要使用完全限定名(例如org.apache.spark.sql.parquet),但对于Spark SQL的内置数据源,也可以使用它们的缩写名(json,parquet,jdbc,orc,libsvm,csv,text)。例如,手动指定csv格式的数据源:

val peopleDFCsv=spark.read.format("csv").load("hdfs://centos01:9000/people.csv")

在指定数据源的同时,可以使用option()方法向指定的数据源传递所需参数。例如,向JDBC数据源传递账号、密码等参数:

val jdbcDF = spark.read.format("jdbc")
  .option("url", "jdbc:mysql://192.168.1.69:3306/spark_db")
  .option("driver","com.mysql.jdbc.Driver")
  .option("dbtable", "student")
  .option("user", "root")
  .option("password", "123456")
  .load()

3、数据写入模式
在写入数据的同时,可以使用mode()方法指定如何处理已经存在的数据,该方法的参数是一个枚举类SaveMode,其取值解析如下:

  • SaveMode.ErrorIfExists:默认值。当向数据源写入一个DataFrame时,如果数据已经存在,则会抛出异常。
  • SaveMode.Append:当向数据源写入一个DataFrame时,如果数据或表已经存在,则会在原有的基础上进行追加。
  • SaveMode.Overwrite:当向数据源写入一个DataFrame时,如果数据或表已经存在,则会将其覆盖(包括数据或表的Schema)。
  • SaveMode.Ignore:当向数据源写入一个DataFrame时,如果数据或表已经存在,则不会写入内容,类似SQL中的“CREATE TABLE IF NOT EXISTS”。

例如,HDFS中有一个JSON格式的文件/people.json,内容:

{"name":"Michael"}
{"name":"Andy", "age":30}
{"name":"Justin", "age":19}

现需要查询该文件中的name列,并将结果写入HDFS的/result目录中,若该目录存在则将其覆盖,代码:

val peopleDF = spark.read.format("json").load("hdfs://centos01:9000/people.json")
peopleDF.select("name")
  .write.mode(SaveMode.Overwrite).format("json")
  .save("hdfs://centos01:9000/result")

4、分区自动推断
表分区是Hive等系统中常用的优化查询效率的方法(Spark SQL的表分区与Hive的表分区类似)。在分区表中,数据通常存储在不同的分区目录中,分区目录通常以“分区列名=值”的格式进行命名。例如,以people作为表名,gendercountry作为分区列,存储数据的目录结构如下:

path
└── to
    └── people
        ├── gender=male
        │   ├── ...
        │   │
        │   ├── country=US
        │   │   └── data.parquet
        │   ├── country=CN
        │   │   └── data.parquet
        │   └── ...
        └── gender=female
            ├── ...
            │
            ├── country=US
            │   └── data.parquet
            ├── country=CN
            │   └── data.parquet
            └── ...

对于所有内置的数据源(包括Text/CSV/JSON/ORC/Parquet),Spark SQL都能够根据目录名自动发现和推断分区信息。分区示例:

  1. 在本地(或HDFS)新建以下三个目录及文件,其中的目录people代表表名,gendercountry代表分区列,people.json存储实际人口数据:

    D:\people\gender=male\country=CN\people.json
    D:\people\gender=male\country=US\people.json
    D:\people\gender=female\country=CN\people.json

    三个people.json文件的数据分别如下:

    {"name":"zhangsan","age":32}
    {"name":"lisi", "age":30}
    {"name":"wangwu", "age":19}
    {"name":"Michael"}
    {"name":"Jack", "age":20}
    {"name":"Justin", "age":18}
    {"name":"xiaohong","age":17}
    {"name":"xiaohua", "age":22}
    {"name":"huanhuan", "age":16}
  2. 执行以下代码,读取表people的数据并显示:

    val usersDF = spark.read.format("json").load("D:\\people") //读取表数据为一个DataFrame
    usersDF.printSchema() //输出Schema信息
    usersDF.show() //输出表数据

    控制台输出的Schema信息如下:

    root
    |-- age: long (nullable = true)
    |-- name: string (nullable = true)
    |-- gender: string (nullable = true)
    |-- country: string (nullable = true)

    控制台输出的表数据如下:

    +----+--------+------+-------+
    | age|    name|gender|country|
    +----+--------+------+-------+
    |  17|xiaohong|female|     CN|
    |  22| xiaohua|female|     CN|
    |  16|huanhuan|female|     CN|
    |  32|zhangsan|  male|     CN|
    |  30|     lisi|  male|     CN|
    |  19|   wangwu|  male|     CN|
    |null| Michael|  male|     US|
    |  20|     Jack|  male|     US|
    |  18|   Justin|  male|     US|
    +----+--------+------+-------+

从控制台输出的Schema信息和表数据可以看出,Spark SQL在读取数据时,自动推断出了两个分区列gendercountry,并将该两列的值添加到了DataFrame中。

Parquet文件

Apache Parquet是Hadoop生态系统中任何项目都可以使用的列式存储格式,不受数据处理框架、数据模型和编程语言的影响。Spark SQL支持对Parquet文件的读写,并且可以自动保存源数据的Schema。当写入Parquet文件时,为了提高兼容性,所有列都会自动转换为“可为空”状态。
加载和写入Parquet文件时,除了可以使用load()方法和save()方法外,还可以直接使用Spark SQL内置的parquet()方法,例如以下代码:

//读取Parquet文件为一个DataFrame
val usersDF = spark.read.parquet("hdfs://centos01:9000/users.parquet")
//将DataFrame相关数据保存为Parquet文件,包括Schema信息
usersDF.select("name","favorite_color")
  .write.parquet("hdfs://centos01:9000/result")

JSON数据集

Spark SQL可以自动推断JSON文件的Schema,并将其加载为DataFrame。在加载和写入JSON文件时,除了可以使用load()方法和save()方法外,还可以直接使用Spark SQL内置的json()方法。该方法不仅可以读写JSON文件,还可以将Dataset[String]类型的数据集转为DataFrame。

需要注意的是,要想成功的将一个JSON文件加载为DataFrame,JSON文件的每一行必须包含一个独立有效的JSON对象,而不能将一个JSON对象分散在多行。例如以下JSON内容可以被成功加载:

{"name":"zhangsan","age":32}
{"name":"lisi", "age":30}
{"name":"wangwu", "age":19}

使用json()方法加载JSON数据的例子如下代码所示:

//创建或得到SparkSession
val spark = SparkSession.builder()
.appName("SparkSQLDataSource")
.config("spark.sql.parquet.mergeSchema",true)
.master("local[*]")
.getOrCreate()

/****1. 创建用户基本信息表*****/
import spark.implicits._
//创建用户信息Dataset集合
val arr=Array(
 "{'name':'zhangsan','age':20}",
 "{'name':'lisi','age':18}"
)
val userInfo: Dataset[String] = spark.createDataset(arr)
//将Dataset[String]转为DataFrame
val userInfoDF = spark.read.json(userInfo)
//创建临时视图user_info
userInfoDF.createTempView("user_info")
//显示数据
userInfoDF.show()
// +---+--------+
// |age|    name|
// +---+--------+
// | 20|zhangsan|
// | 18|    lisi|
// +---+--------+

/****2. 创建用户成绩表*****/
//读取JSON文件
val userScoreDF = spark.read.json("D:\\people\\people.json")
//创建临时视图user_score
userScoreDF.createTempView("user_score")
userScoreDF.show()
// +--------+-----+
// |    name|score|
// +--------+-----+
// |zhangsan|   98|
// |    lisi|   88|
// |  wangwu|   95|
//      +--------+-----+
/****3. 根据name字段关联查询*****/
val resDF=spark.sql("SELECT i.age,i.name,c.score FROM user_info i " +
                      "JOIN user_score c ON i.name=c.name") 
resDF.show()
// +---+--------+-----+
// |age|    name|score|
// +---+--------+-----+
// | 20|zhangsan|   98|
// | 18|    lisi|   88|
// +---+--------+-----+

Hive 表

Spark SQL还支持读取和写入存储在Apache Hive中的数据。然而,由于Hive有大量依赖项,这些依赖项不包括在默认的Spark发行版中,如果在classpath上配置了这些Hive依赖项,Spark将自动加载它们。需要注意的是,这些Hive依赖项必须出现在所有Worker节点上,因为它们需要访问Hive序列化和反序列化库(SerDes),以便访问存储在Hive中的数据。
使用Spark SQL读取和写入Hive数据:
1、创建SparkSession对象
创建一个SparkSession对象,并开启Hive支持,代码:

val spark = SparkSession
  .builder()
  .appName("Spark Hive Demo")
  .enableHiveSupport()//开启Hive支持
  .getOrCreate()

2、创建Hive表
创建一张Hive表students,并指定字段分隔符为制表符“\t”,代码:

spark.sql("CREATE TABLE IF NOT EXISTS students (name STRING, age INT) " +
  "ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'")

3、导入本地数据到Hive表
本地文件/home/hadoop/students.txt的内容如下(字段之间以制表符“\t”分隔):

zhangsan    20
lisi    25
wangwu  19

将本地文件/home/hadoop/students.txt中的数据导入到表students中,代码:

spark.sql("LOAD DATA LOCAL INPATH '/home/hadoop/students.txt'  INTO TABLE students")

4、查询表数据
查询表students的数据并显示到控制台,代码:

spark.sql("SELECT * FROM students").show()

显示结果:

+--------+---+
|    name|age|
+--------+---+
|zhangsan| 20|
|     lisi| 25|
|   wangwu| 19|
+--------+---+

5、创建表的同时指定存储格式
创建一个Hive表hive_records,数据存储格式为Parquet(默认为普通文本格式),代码:

spark.sql("CREATE TABLE hive_records(key STRING, value INT) STORED AS PARQUET")

6、将DataFrame写入Hive表
使用saveAsTable()方法可以将一个DataFrame写入到指定的Hive表中。例如,加载students表的数据并转为DataFrame,然后将DataFrame写入Hive表hive_records中,代码:

//加载students表的数据为DataFrame
val studentsDF = spark.table("students")
//将DataFrame以覆盖的方式写入表hive_records中
studentsDF.write.mode(SaveMode.Overwrite).saveAsTable("hive_records")
//查询hive_records表数据并显示到控制
spark.sql("SELECT * FROM hive_records").show()

Spark SQL应用程序写完后,需要提交到Spark集群中运行。若以Hive为数据源,提交之前需要做好Hive数据仓库、元数据库等的配置。

JDBC

Spark SQL还可以使用JDBC API从其他关系型数据库读取数据,返回的结果仍然是一个DataFrame,可以很容易地在Spark SQL中处理,或者与其他数据源进行连接查询。在使用JDBC连接数据库时可以指定相应的连接属性,常用的连接属性如表。

file

使用JDBC API对MySQL表student和表score进行关联查询,代码:

val jdbcDF = spark.read.format("jdbc")
  .option("url", "jdbc:mysql://192.168.1.69:3306/spark_db")
  .option("driver","com.mysql.jdbc.Driver")
  .option("dbtable", "(select st.name,sc.score from student st,score sc " +
    "where st.id=sc.id) t")
  .option("user", "root")
  .option("password", "123456")
  .load()

上述代码中,dbtable属性的值是一个子查询,相当于SQL查询中的FROM关键字后的一部分。除了上述查询方式外,使用query属性编写完整SQL语句进行查询也能达到同样的效果,代码:

val jdbcDF = spark.read.format("jdbc")
  .option("url", "jdbc:mysql://192.168.1.234:3306/spark_db")
  .option("driver","com.mysql.jdbc.Driver")
  .option("query", "select st.name,sc.score from student st,score sc " +
    "where st.id=sc.id")
  .option("user", "root")
  .option("password", "123456")
  .load()

Spark SQL内置函数

Spark SQL内置了大量的函数,位于API org.apache.spark.sql.functions中。这些函数主要分为10类:UDF函数、聚合函数、日期函数、排序函数、非聚合函数、数学函数、混杂函数、窗口函数、字符串函数、集合函数,大部分函数与Hive中相同。
使用内置函数有两种方式:一种是通过编程的方式使用;另一种是在SQL语句中使用。例如,以编程的方式使用lower()函数将用户姓名转为小写,代码如下:

//显示DataFrame数据(df指DataFrame对象)
df.show()
// +--------+
// |    name|
// +--------+
// |ZhangSan|
// |    LiSi|
// |  WangWu|
// +--------+
//使用lower()函数将某列转为小写
import org.apache.spark.sql.functions._
df.select(lower(col("name")).as("name")).show()
// +--------+
// |    name|
// +--------+
// |zhangsan|
// |    lisi|
// |  wangwu|
// +--------+

Spark SQL自定义函数

当Spark SQL提供的内置函数不能满足查询需求时,用户也可以根据自己的业务编写自定义函数
(User Defined Functions,UDF),然后在Spark SQL中调用。
Spark SQL提供了一些常用的聚合函数,如count()countDistinct()avg()max()min()等。此外,用户也可以根据自己的业务编写自定义聚合函数(User Defined Aggregate Functions,UDAF)。
UDF主要是针对单个输入,返回单个输出;而UDAF则可以针对多个输入进行聚合计算返回单个输出,功能更加强大。要编写UDAF,需要新建一个类,继承抽象类UserDefinedAggregateFunction,并实现其中未实现的方法。

Spark SQL开窗函数

row_number()开窗函数是Spark SQL中常用的一个窗口函数,使用该函数可以在查询结果中对每个分组的数据,按照其排序的顺序添加一列行号(从1开始),根据行号可以方便的对每一组数据取前N行(分组取TOPN)。row_number()函数的使用格式如下:

row_number() over (partition by 列名 order by 列名 desc) 行号列别名

格式说明:

  • partition by:按照某一列进行分组。
  • order by:分组后按照某一列进行组内排序。
  • desc:降序,默认升序。

Views: 148

Maven打包的三种方式

Maven可以使用mvn package指令对项目进行打包,如果使用Java -jar xxx.jar执行运行jar文件,会出现"no main manifest attribute, in xxx.jar"(没有设置Main-Class)、ClassNotFoundException(找不到依赖包)等错误。

要想jar包能直接通过java -jar xxx.jar运行,需要满足:

  1. 在jar包中的META-INF/MANIFEST.MF中指定Main-Class属性,这样才能确定程序的入口在哪里;

  2. 要能加载到依赖包。

使用Maven有以下几种方法可以生成能直接运行的jar包,可以根据需要选择一种合适的方法。

方法一:使用maven-jar-plugin和maven-dependency-plugin插件打包

在pom.xml中配置:

<build>  
    <plugins>  
        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-jar-plugin</artifactId>  
            <version>2.6</version>  
            <configuration>  
                <archive>  
                    <manifest>  
                        <addClasspath>true</addClasspath>  
                        <classpathPrefix>lib/</classpathPrefix>  
                        <mainClass>com.xxx.Main</mainClass>  
                    </manifest>  
                </archive>  
            </configuration>  
        </plugin>  
        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-dependency-plugin</artifactId>  
            <version>2.10</version>  
            <executions>  
                <execution>  
                    <id>copy-dependencies</id>  
                    <phase>package</phase>  
                    <goals>  
                        <goal>copy-dependencies</goal>  
                    </goals>  
                    <configuration>  
                        <outputDirectory>${project.build.directory}/lib</outputDirectory>  
                    </configuration>  
                </execution>  
            </executions>  
        </plugin>  
    </plugins>  
</build>  

maven-jar-plugin用于生成META-INF/MANIFEST.MF文件的部分内容,<mainClass>com.xxx.Main</mainClass>指定MANIFEST.MF中的Main-Class属性值,<addClasspath>true</addClasspath>会在MANIFEST.MF加上Class-Path项并配置依赖包,<classpathPrefix>lib/</classpathPrefix>指定依赖包所在目录。

例如下面是一个通过maven-jar-plugin插件生成的MANIFEST.MF文件片段:

Class-Path: lib/commons-logging-1.2.jar lib/commons-io-2.4.jar  
Main-Class: com.xxx.Main  

只是生成MANIFEST.MF文件还不够,maven-dependency-plugin插件用于将依赖包拷贝到<outputDirectory>${project.build.directory}/lib</outputDirectory>指定的位置,即lib目录下。

配置完成后,通过mvn package指令打包,会在target目录下生成jar包,并将依赖包拷贝到target/lib目录下,目录结构如下:

file

指定了Main-Class,有了依赖包,那么就可以直接通过java -jar xxx.jar运行jar包。

这种方式生成jar包有个缺点,就是生成的jar包太多不便于管理,下面两种方式只生成一个jar文件,包含项目本身的代码、资源以及所有的依赖包。

方法二:使用maven-assembly-plugin插件打包

在pom.xml中配置:

<build>  
    <plugins>  

        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-assembly-plugin</artifactId>  
            <version>2.5.5</version>  
            <configuration>  
                <archive>  
                    <manifest>  
                        <mainClass>com.xxx.Main</mainClass>  
                    </manifest>  
                </archive>  
                <descriptorRefs>  
                    <descriptorRef>jar-with-dependencies</descriptorRef>  
                </descriptorRefs>  
            </configuration>  
        </plugin>  

    </plugins>  
</build>

打包方式:

mvn package assembly:single

打包后会在target目录下生成一个xxx-jar-with-dependencies.jar文件,这个文件不但包含了自己项目中的代码和资源,还包含了所有依赖包的内容。所以可以直接通过java -jar来运行。

此外还可以直接通过mvn package来打包,无需assembly:single,不过需要加上一些配置:

<build>  
    <plugins>  

        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-assembly-plugin</artifactId>  
            <version>2.5.5</version>  
            <configuration>  
                <archive>  
                    <manifest>  
                        <mainClass>com.xxx.Main</mainClass>  
                    </manifest>  
                </archive>  
                <descriptorRefs>  
                    <descriptorRef>jar-with-dependencies</descriptorRef>  
                </descriptorRefs>  
            </configuration>  
            <executions>  
                <execution>  
                    <id>make-assembly</id>  
                    <phase>package</phase>  
                    <goals>  
                        <goal>single</goal>  
                    </goals>  
                </execution>  
            </executions>  
        </plugin>  

    </plugins>  
</build>  

其中<phase>package</phase>、<goal>single</goal>即表示在执行package打包时,执行assembly:single,所以可以直接使用mvn package打包。
不过,如果项目中用到spring Framework,用这种方式打出来的包运行时会出错,使用下面的方法三可以处理。

方法三:使用maven-shade-plugin插件打包
在pom.xml中配置:

<build>  
    <plugins>  

        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-shade-plugin</artifactId>  
            <version>2.4.1</version>  
            <executions>  
                <execution>  
                    <phase>package</phase>  
                    <goals>  
                        <goal>shade</goal>  
                    </goals>  
                    <configuration>  
                        <transformers>  
                            <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">  
                                <mainClass>com.xxx.Main</mainClass>  
                            </transformer>  
                        </transformers>  
                    </configuration>  
                </execution>  
            </executions>  
        </plugin>  

    </plugins>  
</build>  

配置完成后,执行mvn package即可打包。在target目录下会生成两个jar包,注意不是original-xxx.jar文件,而是另外一个。和maven-assembly-plugin一样,生成的jar文件包含了所有依赖,所以可以直接运行。
如果项目中用到了Spring Framework,将依赖打到一个jar包中,运行时会出现读取XML schema文件出错。原因是Spring Framework的多个jar包中包含相同的文件spring.handlersspring.schemas,如果生成一个jar包会互相覆盖。为了避免互相影响,可以使用AppendingTransformer来对文件内容追加合并:

<build>  
    <plugins>  
        <plugin>  
            <groupId>org.apache.maven.plugins</groupId>  
            <artifactId>maven-shade-plugin</artifactId>  
            <version>2.4.1</version>  
            <executions>  
                <execution>  
                    <phase>package</phase>  
                    <goals>  
                        <goal>shade</goal>  
                    </goals>  
                    <configuration>  
                        <transformers>  
                            <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">  
                                <mainClass>com.xxx.Main</mainClass>  
                            </transformer>  
                            <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">  
                                <resource>META-INF/spring.handlers</resource>  
                            </transformer>  
                            <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">  
                                <resource>META-INF/spring.schemas</resource>  
                            </transformer>  
                        </transformers>  
                    </configuration>  
                </execution>  
            </executions>  
        </plugin>  
    </plugins>  
</build>  

配置完成后,执行mvn package即可打包。

Views: 134

配置 Flume Source

安装netcat

Netcat 是一款简单的Unix工具,简称 nc,安全界叫它瑞士军刀, 使用UDP和TCP协议。 它是一个可靠的容易被其他程序所启用的后台操作工具,同时它也被用作网络的测试工具或黑客工具。 使用它你可以轻易的建立任何连接。内建有很多实用的工具。

$> sudo yum install nmap-ncat.x86_64

nc的一些用法:

端口测试

检测主机上8080端口服务是否开放

telnet 192.168.1.2 8080

或者

nc -vz 192.168.1.2 8080

z表示不发送数据,v表示显示额外信息

nc 命令后面的 8080 可以写成一个范围进行扫描:

nc -v -v -w3 -z 192.168.1.2 8080-8083

两次 -v 是让它报告更详细的内容,-w3 是设置扫描超时时间为 3 秒。

传输测试

A 主机上监听了 8080 端口

nc -l -p 8080

然后在 B 主机上连接过去:

nc 192.168.1.2 8080

两边就可以会话了,随便输入点什么按回车,另外一边应该会显示出来.

netcat source

配置

[hello.conf]

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

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

#定义sink信息
a1.sinks.k1.type=logger

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

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

运行

  1. 启动flume agent$> bin/flume-ng agent -f conf/hello.conf -n a1 -Dflume.root.logger=INFO,console
  2. 启动nc的客户端$> nc localhost 8888 $nc> hello world
  3. 在Flume的终端输出hello world.

exec source

实时日志收集,实时收集日志。

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

a1.sources.r1.type=exec
a1.sources.r1.command=tail -F /home/centos/test.txt

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

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

spooldir 源

监控一个文件夹,静态文件, 批量收集。
收集完之后,会重命名文件成新文件.COMPLETED.

配置文件

[spooldir_r.conf]

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

a1.sources.r1.type=spooldir
a1.sources.r1.spoolDir=/home/centos/spool
a1.sources.r1.fileHeader=true

a1.sinks.k1.type=logger

a1.channels.c1.type=memory

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

创建目录

$>mkdir ~/spool

启动flume

$>bin/flume-ng agent -f ../conf/helloworld.conf -n a1 -Dflume.root.logger=INFO,console

seq source

生成事件序列的源, 一般用于测试.
[seq]

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

a1.sources.r1.type=seq
a1.sources.r1.totalEvents=1000

a1.sinks.k1.type=logger

a1.channels.c1.type=memory

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

[运行]

$>bin/flume-ng agent -f ../conf/helloworld.conf -n a1 -Dflume.root.logger=INFO,console

Stress Source

用于压力测试的源.

a1.sources = stresssource-1
a1.channels = memoryChannel-1
a1.sources.stresssource-1.type = org.apache.flume.source.StressSource
a1.sources.stresssource-1.size = 10240
a1.sources.stresssource-1.maxTotalEvents = 1000000
a1.sources.stresssource-1.channels = memoryChannel-1

TailDir Source

Taildir Source目前只是个预览版本,还不能运行在windows系统上。

Taildir Source监控指定的一些文件,并在检测到新的一行数据产生的时候几乎实时地读取它们,如果新的一行数据还没写完,Taildir Source会等到这行写完后再读取。

Taildir Source是可靠的,即使发生文件滚动也不会丢失数据。它会定期地以JSON格式在一个专门用于定位的文件上记录每个文件的最后读取位置。如果Flume由于某种原因停止或挂掉,它可以从文件的标记位置重新开始读取。

Taildir Source还可以从任意指定的位置开始读取文件。默认情况下,它将从每个文件的第一行开始读取。

文件按照修改时间的顺序来读取。修改时间最早的文件将最先被读取(简单记成:先来先走)。

Taildir Source不重命名、删除或修改它监控的文件。当前不支持读取二进制文件。只能逐行读取文本文件。

文件滚动(file rotate)就是我们常见的log4j等日志框架或者系统会自动丢弃日志文件中时间久远的日志,一般按照日志文件大小或时间来自动分割或丢弃的机制。

练习

使用Flume监听整个目录的实时追加文件,并上传至HDFS

#步骤一:agent Name
a1.sources = r1
a1.sinks = k1
a1.channels = c1

#步骤二:source
# Describe/configure the source
a1.sources.r1.type = TAILDIR 
a1.sources.r1.positionFile = /opt/module/flume/tail_dir.json -- 指定position_file 的位置(记录每次上传后的偏移量,实现断点续传的关键)
a1.sources.r1.filegroups = f1 f2 -- 监控的文件目录集合
a1.sources.r1.filegroups.f1 = /opt/module/flume/files/.*file.* -- 定义监控的文件目录1
a1.sources.r1.filegroups.f2 = /opt/module/flume/files/.*log.* -- 定义监控的文件目录2

#步骤三: channel selector
a1.sources.r1.selector.type = replicating

#步骤四: channel
# Describe the channel
a1.channels.c1.type = memory
a1.channels.c1.capacity = 1000
a1.channels.c1.transactionCapacity = 100

#步骤五: sinkprocessor,默认配置defaultsinkprocessor
#步骤六: sink
a1.sinks.k1.type = hdfs
a1.sinks.k1.hdfs.path = hdfs://hadoop102:9820/flume/upload3/%Y%m%d/%H

#上传文件的前缀
a1.sinks.k1.hdfs.filePrefix = upload-
#是否按照时间滚动文件夹
a1.sinks.k1.hdfs.round = true
#多少时间单位创建一个新的文件夹
a1.sinks.k1.hdfs.roundValue = 1
#重新定义时间单位
a1.sinks.k1.hdfs.roundUnit = hour
#是否使用本地时间戳
a1.sinks.k1.hdfs.useLocalTimeStamp = true
#积攒多少个Event才flush到HDFS一次
a1.sinks.k1.hdfs.batchSize = 100
#设置文件类型,可支持压缩
a1.sinks.k1.hdfs.fileType = DataStream
#多久生成一个新的文件
a1.sinks.k1.hdfs.rollInterval = 60 
#设置每个文件的滚动大小大概是128M
a1.sinks.k1.hdfs.rollSize = 134217700
#文件的滚动与Event数量无关
a1.sinks.k1.hdfs.rollCount = 0
#步骤七:连接source、channel、sink
a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

Views: 342

Sqoop环境安装

问题陈述

Sqoop是Hadoop生态系统和RDBMS之间进行数据传输的一个工具。在学习Sqoop之前首先需要完成学习环境的搭建。这里为了学习方便,采用单机部署方式。

最初Sqoop是Hadoop的一个子项目,它设计只能在Linux操作系统上运行。

先决条件

安装Sqoop的必要前提条件是:

  • 准备Linux操作系统(Centos7)
  • 安装Java环境(JDK1.8)
  • 安装Hadoop环境(Hadoop 3.1.4)

另外为了学习Sqoop的大部分功能,还需要需要安装:

  • Zookeeper 3.5.9)
  • HBase 2.2.3
  • MySQL 5.7
  • Hive 3.1.2

解决办法

准备Linux操作系统

配置主机名映射

设置网络:

虚拟机安装Centos7时
硬盘设置大一些,如40G或更多
不要设置预先分配磁盘空间.
网络适配器设置为NAT连接

打开虚拟网络编辑器

记住NAT模式所在的虚拟网卡对应的子网IP,子网掩码以及网关IP.

然后进入虚拟机终端,设置静态IP

$ vi /etc/sysconfig/network-scripts/ifcfg-ens33
NAME="ens33"
TYPE="Ethernet"
DEVICE="ens33"
BROWSER_ONLY="no"
DEFROUTE="yes"
PROXY_METHOD="none"
IPV4_FAILURE_FATAL="no"
IPV6INIT="yes"
IPV6_AUTOCONF="yes"
IPV6_DEFROUTE="yes"
IPV6_FAILURE_FATAL="no"
IPV6_ADDR_GEN_MODE="stable-privacy"
IPV6_PRIVACY="no"
UUID="b00b9ac0-60c2-4d34-ab88-2413055463cf"

ONBOOT="yes"
BOOTPROTO="static"
IPADDR="192.168.186.100"
PREFIX="24"
GATEWAY="192.168.186.2"
DNS1="223.5.5.5"
DNS2="8.8.8.8"

其中需要修改的是:

  • ONTBOOT设置yes可以实现自动联网
  • BOOTPROTO="static" 设置静态IP,防止IP发生变化
  • IPADDR的前三段要和NAT虚拟网卡的子网IP一致,且第四段在0~254之间选择,又不能和NAT虚拟网卡的子网掩码和其他相同网络中的主机IP重复.
  • PREFIX=24是设置子网掩码的位数长度,换算十进制就是255.255.255.0,因此PREFIX=24也可以直接替换成NETMASK="255.255.255.0"
  • DNS1设置的是阿里的公共DNS地址"223.5.5.5",DNS2设置的是谷歌的公共DNS地址"8,8,8,8"

设置好了之后需要重新启动网络:

$ sudo service network restart
Restarting network (via systemctl):                        [  确定  ]
$ sudo service network status
已配置设备:
lo ens33
当前活跃设备:
lo ens33

查看本机ip:

$ ip a
1: lo: <LOOPBACK,UP,LOWER_UP> mtu 65536 qdisc noqueue state UNKNOWN group default qlen 1000
  link/loopback 00:00:00:00:00:00 brd 00:00:00:00:00:00
  inet 127.0.0.1/8 scope host lo
    valid_lft forever preferred_lft forever
  inet6 ::1/128 scope host
    valid_lft forever preferred_lft forever
2: ens33: <BROADCAST,MULTICAST,UP,LOWER_UP> mtu 1500 qdisc pfifo_fast state UP group default qlen 1000
  link/ether 00:0c:29:b2:5b:75 brd ff:ff:ff:ff:ff:ff
  inet 192.168.186.100/24 brd 192.168.186.255 scope global noprefixroute ens33
    valid_lft forever preferred_lft forever
  inet6 fe80::40db:ee8d:77c1:fbf7/64 scope link noprefixroute
    valid_lft forever preferred_lft forever

设置静态主机名:

# hostnamectl --static set-hostname hadoop100

在hosts文件中配置主机名和本机ip之间的映射关系
(注释掉localhost的部分)

vi /etc/hosts
#127.0.0.1  localhost localhost.localdomain localhost4 localhost4.localdomain4
#::1     localhost localhost.localdomain localhost6 localhost6.localdomain6
192.168.186.100 hadoop100

创建专用账号

不建议直接使用root账号, 这里我们创建一个hadoop账号用于接下来的所有操作.

创建hadoop用户

# useradd hadoop    

设置hadoop用户的密码

# password hadoop

为hadoop账号设置sudoer权限

vi /etc/sudoers

找到root ALL=(ALL) ALL在下面添加一行

## Allow root to run any commands anywhere
root    ALL=(ALL)   ALL
delucia ALL=(ALL)   NOPASSWD:ALL

修改完成后,切换到Hadoop用户

# su hadoop

以后在使用需要root权限的命令时,就可以在命令前面加上sudo来提升权限,且无需输入hadoop密码, 如:

$ sudo ls /root

接下来的所有操作, 如果没有特殊说明, 一律使用hadoop账号.

配置SSH免密登录

由于Hadoop集群的机器之间ssh通信默认需要输入密码,在集群运行时我们不可能为每一次通信都手动输入密码,因此需要配置机器之间的ssh的免密登录。单机伪分布式的Hadoop环境样需要配置本地对本地ssh连接的免密,流程如下:

  1. 首先ssh-keygen命令生成RSA加密的密钥对(公钥和私钥)。

    $ ssh-keygen -t rsa
    Generating public/private rsa key pair.
    Enter file in which to save the key (/home/hadoop/.ssh/id_rsa): 
    Created directory '/home/hadoop/.ssh'.
    Enter passphrase (empty for no passphrase): 
    Enter same passphrase again: 
    Your identification has been saved in /home/hadoop/.ssh/id_rsa.
    Your public key has been saved in /home/hadoop/.ssh/id_rsa.pub.
    The key fingerprint is:
    SHA256:yQYChs4eniVLeICaI2bCB9HopbXUBE9v0lpLjBACUnM hadoop@hadoop100
    The key's randomart image is:
    +---[RSA 2048]----+
    |==O+Eo      |
    |=+.O+.=     |
    |Bo* o+.B     |
    |B%.+ .*o..    |
    |OoB . .S    |
    | =   .     |
    |         |
    +----[SHA256]-----+
  2. 将生成的公钥添加到~/.ssh目录下的authorized_keys文件中。并为authorized_keys文件设置600权限。

    $ cd ~/.ssh/
    $ cat id_rsa.pub >> authorized_keys
    $ chmod 600 authorized_keys

    以上三个命令可以使用一个命令代替:

    $ ssh-copy-id hadoop100
  3. 使用ssh命令连接本地终端,如果不需要输入密码则说明本地的SSH免密配置成功。

    $ ssh hadoop@hadoop100
    The authenticity of host 'hadoop100 (192.168.186.100)' can't be established.
    ECDSA key fingerprint is SHA256:aGLhdt3bIuqtPgrFWnhgrfTKUbDh4CWVTfIgr5E5oV0.
    ECDSA key fingerprint is MD5:b8:bd:b3:65:fe:77:2c:06:2d:ec:58:3a:97:51:dd:ca.
    Are you sure you want to continue connecting (yes/no)? yes
    Warning: Permanently added 'hadoop100,192.168.186.100' (ECDSA) to the list of known hosts.
    Last login: Sat Jan 9 10:16:53 2021 from 192.168.186.1
    
  4. 登出

    $ exit
    Connection to hadoop100 closed.

配置时间同步

集群中的通信和文件传输一般是以系统时间作为约定条件的。所以当集群中机器之间系统如果不一致可能导致各种问题发生,比如访问时间过长,甚至失败。所以配置机器之间的时间同步非常重要。不过由于我们使用的学习环境是单机部署,所以无需配置时间同步。

统一目录结构

目录规划如下:

/opt/ 
  ├── bin         # 脚本和命令
  ├── data        # 程序需要使用的数据
  ├── download    # 下载的软件安装包
  ├── pkg         # 解压方式安装的软件
  └── tmp         # 存放程序生成的临时文件

使用hadoop账户创建目录:

$ sudo mkdir /opt/download
$ sudo mkdir /opt/data
$ sudo mkdir /opt/bin
$ sudo mkdir /opt/tmp
$ sudo mkdir /opt/pkg

为了使用方便更改opt下目录的用户及其所在用户组为hadoop:

$ sudo chown hadoop:hadoop /opt/*
$ ls
bin data download pkg tmp

$ ll
总用量 0
drwxr-xr-x. 2 hadoop hadoop 6 1月  9 14:29 bin
drwxr-xr-x. 2 hadoop hadoop 6 1月  9 14:29 data
drwxr-xr-x. 2 hadoop hadoop 6 1月  9 14:29 download
drwxr-xr-x. 2 hadoop hadoop 6 1月  9 14:37 pkg
drwxr-xr-x. 2 hadoop hadoop 6 1月  9 14:36 tmp

安装Java环境

检查是否已经装过Java JDK

$ rpm -qa | grep java

$ yum list installed | grep java

如果没有安装过Java 则需要安装,JDK版本建议1.8+

Java JDK的下载和解压

去官网下载JDK1.8的安装包,上传到/opt/download/,然后解压到/opt/pkg

$ tar -zxvf jdk-8u261-linux-x64.tar.gz 
$ mv jdk1.8.0_261 /opt/pkg/java

配置java环境变量

确认当前目录为jdk的解压路径

$ pwd
/opt/pkg/java

编辑/etc/profile.d/hadoop.env.sh配置文件(没有则创建)

$ sudo vim /etc/profile.d/hadoop.env.sh

添加新的环境变量配置

# JAVA_HOME
export JAVA_HOME=/opt/pkg/java
PATH=$JAVA_HOME/bin:$PATH

export PATH

使新的环境变量立刻生效

$ source /etc/profile.d/env.sh

验证环境变量

$ java -version
$ java
$ javac

安装Hadoop

下载和解压

从官网下载hadoop-3.1.4.tar.gz上传到服务器,解压到指定目录:

$ tar -zxvf hadoop.tar.gz -C /opt/pkg/

编辑/etc/profile.d/env.sh配置文件,添加环境变量:

# HADOOP_HOME
export HADOOP_HOME=/opt/pkg/hadoop
export PATH=$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$PATH

使新的环境变量立刻生效:

$ source /etc/profile.d/env.sh

验证:

$ hadoop version

修改Hadoop相关命令执行环境

找到Hadoop安装目录下的hadoop/etc/hadoop/hadoop-env.sh文件,找到这一处将JAVA_HOME修改为真实JDK路径即可:

# The java implementation to use.
export JAVA_HOME=/opt/pkg/java

找到hadoop/etc/hadoop/yarn.env.sh文件,做同样修改:

# export JAVA_HOME=/home/y/libexec/jdk1.6.0/
export JAVA_HOME=/opt/pkg/java

找到hadoop/etc/hadoop/mapred.env.sh,做同样修改:

# export JAVA_HOME=/home/y/libexec/jdk1.6.0/
export JAVA_HOME=/opt/pkg/java

修改Hadoop配置

来到 hadoop/etc/hadoop/,修改以下配置文件。

1) hadoop/etc/hadoop/core-site.xml – Hadoop核心配置文件

    <configuration>
      <!-- 指定NameNode的地址和端口. -->
      <property>
        <name>fs.defaultFS</name>
        <value>hdfs://hadoop100:8020</value>
      </property>

      <!-- 指定HDFS系统运行时产生的文件的存储目录. -->
      <property>
        <name>hadoop.tmp.dir</name>
        <value>/opt/pkg/hadoop/data/tmp</value>
      </property>

     <!-- 缓冲区大小,实际工作中根据服务器性能动态调整;默认值4096 -->
     <property>
       <name>io.file.buffer.size</name>
       <value>4096</value>
     </property>

     <!-- 开启hdfs的垃圾桶机制,删除掉的数据可以从垃圾桶中回收,单位分钟;默认值0 -->
     <property>
       <name>fs.trash.interval</name>
       <value>10080</value>
     </property>
   </configuration>

注意:主机名要修改成本机的实际主机名。

hadoop.tmp.dir十分重要,此目录下保存hadoop集群中namenode和datanode的所有数据。

2) hadoop/etc/hadoop/hdfs-site.xml – HDFS相关配置

    <configuration>
      <!-- 设置HDFS中的数据副本数. -->
      <property>
        <name>dfs.replication</name>
        <value>1</value>
      </property>

      <!-- 设置Hadoop的Secondary NameNode的主机配置 -->
      <property>
        <name>dfs.namenode.secondary.http-address</name>
        <value>hadoop100:9868</value>
      </property>

      <property>
         <name>dfs.namenode.http-address</name>
        <value>hadoop100:9870</value>
      </property>

      <!-- 是否检查操作HDFS文件系统的用户权限. -->
      <property>
      <name>dfs.permissions</name>
      <value>false</value>
     </property>
   </configuration>

dfs.replication默认是3,为了节省虚拟机资源,这里设置为1

​ 全分布式情况下,SecondaryNameNode和NameNode 应分开部署

dfs.namenode.secondary.http-address默认就是本地,如果是伪分布式可以不用配置

3) hadoop/etc/hadoop/mapred-site.xml – mapreduce 相关配置

   <configuration>
     <!-- 指定MapReduce程序由Yarn进行调度. -->
     <property>
       <name>mapreduce.framework.name</name>
       <value>yarn</value>
     </property>

     <!-- Mapreduce的Job历史记录服务器主机端口设置. -->
     <property>
       <name>mapreduce.jobhistory.address</name>
       <value>hadoop100:10020</value>
     </property>

     <!-- Mapreduce的Job历史记录的Webapp端地址. -->
     <property>
       <name>mapreduce.jobhistory.webapp.address</name>
       <value>hadoop100:19888</value>
     </property>

     <property>
       <name>yarn.app.mapreduce.am.env</name>
       <value>HADOOP_MAPRED_HOME=/opt/pkg/hadoop</value>
     </property>

     <property>
       <name>mapreduce.map.env</name>
       <value>HADOOP_MAPRED_HOME=/opt/pkg/hadoop</value>
     </property>

     <property>
       <name>mapreduce.reduce.env</name>
       <value>HADOOP_MAPRED_HOME=/opt/pkg/hadoop</value>
     </property>
   </configuration>

mapreduce.jobhistory相关配置是可选配置,用于查看MR任务的历史日志。

​ 这里主机名千万不要弄错,不然任务执行会失败,且不容易找原因。

​ 需要手动启动MapReduceJobHistory后台服务才能在Yarn的页面打开历史日志。

4) 配置 yarn-site.xml

    <configuration>
     <!-- 设置Yarn的ResourceManager节点主机名. -->
     <property>
       <name>yarn.resourcemanager.hostname</name>
       <value>hadoop100</value>
     </property>

     <!-- 设置Mapper端将数据发送到Reducer端的方式. -->
     <property>
       <name>yarn.nodemanager.aux-services</name>
       <value>mapreduce_shuffle</value>
     </property>

     <!-- 是否开启日志手机功能. -->
     <property>
       <name>yarn.log-aggregation-enable</name>
       <value>true</value>
     </property>

     <!-- 日志保留时间(7天). -->
     <property>
       <name>yarn.log-aggregation.retain-seconds</name>
       <value>604800</value>
     </property>

     <!-- 如果vmem、pmem资源不够,会报错,此处将资源监察置为false -->
     <property>
       <name>yarn.nodemanager.vmem-check-enabled</name>
       <value>false</value>
     </property>

     <property>
       <name>yarn.nodemanager.pmem-check-enabled</name>
       <value>false</value>
     </property>
    </configuration>

5) workers DataNode 节点配置

   vi workers

   $ vi workers
   hadoop100

这里单机伪分布式环境可以不进行修改,默认是localhost, 也可以改成本机的主机名。

全分布式配置则需要每行输入一个DataNode主机名。

注意DataNode的主机名中不要有空格和空行,因为其他脚本会获取相关主机名信息。

格式化名称节点

$ hdfs namenode -format
21/01/09 19:27:21 INFO namenode.NameNode: STARTUP_MSG: 
/************************************************************
STARTUP_MSG: Starting NameNode
STARTUP_MSG:  host = hadoop100/192.168.186.100
STARTUP_MSG:  args = [-format]
STARTUP_MSG:  version = 2.7.3
************************************************************/
21/01/09 19:27:21 INFO namenode.NameNode: registered UNIX signal handlers for [TERM, HUP, INT]
21/01/09 19:27:21 INFO namenode.NameNode: createNameNode [-format]
Formatting using clusterid: CID-08318e9e-e202-48f3-bcb1-548ca50310c9
21/01/09 19:27:22 INFO util.GSet: Computing capacity for map BlocksMap
21/01/09 19:27:22 INFO util.GSet: VM type    = 64-bit
21/01/09 19:27:22 INFO util.GSet: 2.0% max memory 966.7 MB = 19.3 MB
21/01/09 19:27:22 INFO util.GSet: capacity   = 2^21 = 2097152 entries
21/01/09 19:27:22 INFO blockmanagement.BlockManager: dfs.block.access.token.enable=false
21/01/09 19:27:22 INFO blockmanagement.BlockManager: defaultReplication     = 1
21/01/09 19:27:22 INFO blockmanagement.BlockManager: maxReplication       = 512
21/01/09 19:27:22 INFO blockmanagement.BlockManager: minReplication       = 1
21/01/09 19:27:22 INFO blockmanagement.BlockManager: maxReplicationStreams   = 2
21/01/09 19:27:22 INFO blockmanagement.BlockManager: replicationRecheckInterval = 3000
21/01/09 19:27:22 INFO blockmanagement.BlockManager: encryptDataTransfer    = false
21/01/09 19:27:22 INFO blockmanagement.BlockManager: maxNumBlocksToLog     = 1000
21/01/09 19:27:22 INFO namenode.FSNamesystem: fsOwner       = hadoop (auth:SIMPLE)
21/01/09 19:27:22 INFO namenode.FSNamesystem: supergroup     = supergroup
21/01/09 19:27:22 INFO namenode.FSNamesystem: isPermissionEnabled = false
21/01/09 19:27:22 INFO namenode.FSNamesystem: HA Enabled: false
21/01/09 19:27:22 INFO namenode.FSNamesystem: Append Enabled: true
21/01/09 19:27:23 INFO common.Storage: Storage directory /opt/pkg/hadoop/data/tmp/dfs/name has been successfully formatted.
/************************************************************
SHUTDOWN_MSG: Shutting down NameNode at hadoop100/192.168.186.100
************************************************************/

运行和测试

启动Hadoop环境,刚启动Hadoop的HDFS系统后会有几秒的安全模式,安全模式期间无法进行任何数据处理,这也是为什么不建议使用start-all.sh脚本一次性启动DFS进程和Yarn进程,而是先启动dfs后过30秒左右再启动Yarn相关进程。

1) 启动所有DFS进程:

   $ start-dfs.sh
   Starting namenodes on [hadoop100]
   hadoop100: starting namenode, logging to /opt/pkg/hadoop/logs/hadoop-hadoop-namenode-hadoop100.out
   hadoop100: starting datanode, logging to /opt/pkg/hadoop/logs/hadoop-hadoop-datanode-hadoop100.out
   Starting secondary namenodes [hadoop100]
   hadoop100: starting secondarynamenode, logging to /opt/pkg/hadoop/logs/hadoop-hadoop-secondarynamenode-hadoop100.out

2) 启动所有YARN进程:

   $ start-yarn.sh
   starting yarn daemons
   starting resourcemanager, logging to /opt/pkg/hadoop/logs/yarn-hadoop-resourcemanager-hadoop100.out
   hadoop100: starting nodemanager, logging to /opt/pkg/hadoop/logs/yarn-hadoop-nodemanager-hadoop100.out

启动MapReduceJobHistory后台服务 – 用于查看MR执行的历史日志

   $ mr-jobhistory-daemon.sh start historyserver

3) 查看是否相关进程都成功启动

执行jps命令,看看是否会有如下进程:

   $ jps
   14608 NodeManager
   14361 SecondaryNameNode
   14203 DataNode
   14510 ResourceManager
   14079 NameNode

4) 单一进程管理

   # 在主节点上使用以下命令启动 HDFS NameNode: 
   hdfs --daemon start namenode

   # 在主节点上使用以下命令启动 HDFS SecondaryNamenode: 
   hdfs --daemon start secondarynamenode

   # 在从节点上使用以下命令启动 HDFS DataNode: 
   hdfs --daemon start datanode

   # 在主节点上使用以下命令启动 YARN ResourceManager: 
   yarn --daemon start resourcemanager

   # 在从节点上使用以下命令启动 YARN nodemanager: 
   yarn --daemon start nodemanager

以上脚本位于$HADOOP_HOME/sbin/目录下。如果想要停止某个节点上某个进程,只需要把命令中的start 改为stop 即可。

Web界面进行验证

访问http://hadoop100:9870查看HDFS情况

image-20220210022917279

访问http://hadoop100:8088查看YARN情况

image-20220210022937527

测试Hadoop集群

使用官方自带的示例程序测试Hadoop集群

启动DFS和YARN进程,找到测试程序的位置:

$ cd /opt/pkg/hadoop/share/hadoop/mapreduce
$ hadoop jar hadoop-mapreduce-examples-2.7.3.jar wordcount
Usage: wordcount <in> [<in>...] <out>

准备输入文件并上传到HDFS系统

$ cat /opt/data/mapred/input/wc.txt
hadoop hadoop hadoop
hi hi hi hello hadoop
hello world hadoop

$ hadoop fs -mkdir -p /input/wc

$ hadoop fs -put wc.txt /input/wc/
Found 1 items
-rw-r--r--  1 hadoop supergroup     62 2021-01-09 20:15 /input/wc/wc.txt

$ hadoop fs -cat /input/wc/wc.txt
hadoop hadoop hadoop
hi hi hi hello hadoop
hello world hadoop

运行官方示例程序wordcount,并将结果输出到/output/wc之中

$ cd /opt/pkg/hadoop/share/hadoop/mapreduce
$ hadoop jar hadoop-mapreduce-examples-3.1.4.jar wordcount /input/wc/ /output/wc/

控制台输出:

2020-01-23 18:38:45,914 INFO client.RMProxy: Connecting to ResourceManager at hadoop100/192.168.186.100:8032
2020-01-23 18:38:47,204 INFO mapreduce.JobResourceUploader: Disabling Erasure Coding for path: /tmp/hadoop-yarn/staging/hadoop/.staging/job_1642908422458_0001
2020-01-23 18:38:47,988 INFO input.FileInputFormat: Total input files to process : 1
2020-01-23 18:38:49,033 INFO mapreduce.JobSubmitter: number of splits:1
2020-01-23 18:38:49,788 INFO mapreduce.JobSubmitter: Submitting tokens for job: job_1642908422458_0001
2020-01-23 18:38:49,790 INFO mapreduce.JobSubmitter: Executing with tokens: []
2020-01-23 18:38:50,108 INFO conf.Configuration: resource-types.xml not found
2020-01-23 18:38:50,108 INFO resource.ResourceUtils: Unable to find 'resource-types.xml'.
2020-01-23 18:38:50,740 INFO impl.YarnClientImpl: Submitted application application_1642908422458_0001
2020-01-23 18:38:50,796 INFO mapreduce.Job: The url to track the job: http://hadoop100:8088/proxy/application_1642908422458_0001/
2020-01-23 18:38:50,797 INFO mapreduce.Job: Running job: job_1642908422458_0001
2020-01-23 18:39:09,424 INFO mapreduce.Job: Job job_1642908422458_0001 running in uber mode : false
2020-01-23 18:39:09,425 INFO mapreduce.Job: map 0% reduce 0%
2020-01-23 18:39:19,633 INFO mapreduce.Job: map 100% reduce 0%
2020-01-23 18:39:28,781 INFO mapreduce.Job: map 100% reduce 100%
2020-01-23 18:39:30,812 INFO mapreduce.Job: Job job_1642908422458_0001 completed successfully
2020-01-23 18:39:30,973 INFO mapreduce.Job: Counters: 53
    File System Counters
        FILE: Number of bytes read=52
        FILE: Number of bytes written=444077
         FILE: Number of read operations=0
        FILE: Number of large read operations=0
        FILE: Number of write operations=0
        HDFS: Number of bytes read=165
        HDFS: Number of bytes written=30
        HDFS: Number of read operations=8
        HDFS: Number of large read operations=0
        HDFS: Number of write operations=2
     Job Counters
        Launched map tasks=1
        Launched reduce tasks=1
        Data-local map tasks=1
        Total time spent by all maps in occupied slots (ms)=8195
        Total time spent by all reduces in occupied slots (ms)=6333
        Total time spent by all map tasks (ms)=8195
        Total time spent by all reduce tasks (ms)=6333
        Total vcore-milliseconds taken by all map tasks=8195
        Total vcore-milliseconds taken by all reduce tasks=6333
        Total megabyte-milliseconds taken by all map tasks=8391680
        Total megabyte-milliseconds taken by all reduce tasks=6484992
    Map-Reduce Framework
        Map input records=4
         Map output records=11
        Map output bytes=106
        Map output materialized bytes=52
        Input split bytes=102
        Combine input records=11
        Combine output records=4
        Reduce input groups=4
        Reduce shuffle bytes=52
        Reduce input records=4
        Reduce output records=4
        Spilled Records=8
        Shuffled Maps =1
        Failed Shuffles=0
        Merged Map outputs=1
        GC time elapsed (ms)=235
        CPU time spent (ms)=3340
        Physical memory (bytes) snapshot=366059520
        Virtual memory (bytes) snapshot=5470892032
        Total committed heap usage (bytes)=291639296
        Peak Map Physical memory (bytes)=233541632
        Peak Map Virtual memory (bytes)=2732072960
        Peak Reduce Physical memory (bytes)=132517888
        Peak Reduce Virtual memory (bytes)=2738819072
    Shuffle Errors
        BAD_ID=0
        CONNECTION=0
        IO_ERROR=0
        WRONG_LENGTH=0
        WRONG_MAP=0
        WRONG_REDUCE=0
    File Input Format Counters
        Bytes Read=63
    File Output Format Counters
        Bytes Written=30

注意,输入是文件夹,可以指定多个。输出是一个必须不存在的文件夹路径。

查看保存在HDFS上的结果:

$ hadoop fs -ls /output/wc/
Found 2 items
-rw-r--r--  1 hadoop supergroup     0 2021-01-09 20:23 /output/wc/_SUCCESS
-rw-r--r--  1 hadoop supergroup     30 2021-01-09 20:23 /output/wc/part-r-00000

$ hadoop fs -cat /output/wc/part-r-00000
hadoop 5
hello  2
hi 3
world  1

在MR任务执行时,可以通过Yarn的Web UI界面查看进度:

image-20220210022958848

执行完毕以后点击指定MapReduce程序的TrackingUI一栏下的History可以查看历史日志记录

image-20220210023011146

如果跳转页面报404说明没有启动JobHistoryServer服务。

也可以在HDFS的Web界面上查看结果。

位置:Utilities > HDFS browser > /output/wc/ > part-r-00000

image-20220210023023180

关闭集群

$ stop-all.sh
This script is Deprecated. Instead use stop-dfs.sh and stop-yarn.sh
Stopping namenodes on [hadoop100]
hadoop100: stopping namenode
hadoop100: stopping datanode
Stopping secondary namenodes [hadoop100]
hadoop100: stopping secondarynamenode
stopping yarn daemons
stopping resourcemanager
hadoop100: stopping nodemanager
no proxyserver to stop

安装HBase

Hbase有自带的Zookeeper, 为了更好的使用建议使用自己安装的zookeeper环境.

安装Zookeeper

从官网下载apache-zookeeper-3.5.9-bin.tar.gz,安装到/opt/pkg/zookeeper,单机模式部署,配置文件 $ZOOKEEPER_HOME/conf/zoo.cfg

# The number of milliseconds of each tick
tickTime=2000
# The number of ticks that the initial
# synchronization phase can take
initLimit=10
# The number of ticks that can pass between
# sending a request and getting an acknowledgement
syncLimit=5
# the directory where the snapshot is stored.
# do not use /tmp for storage, /tmp here is just
# example sakes.
dataDir=/opt/tmp/zookeeper/data
dataLogDir=/opt/tmp/zookeeper/dataLog
# the port at which the clients will connect
clientPort=2181
# the maximum number of client connections.
# increase this if you need to handle more clients
#maxClientCnxns=60
#
# Be sure to read the maintenance section of the
# administrator guide before turning on autopurge.
#
# http://zookeeper.apache.org/doc/current/zookeeperAdmin.html#sc_maintenance
#
# The number of snapshots to retain in dataDir
#autopurge.snapRetainCount=3
# Purge task interval in hours
# Set to "0" to disable auto purge feature
#autopurge.purgeInterval=1

下载安装包

安装包下载地址:https://www.apache.org/dyn/closer.lua/hbase/2.2.6/hbase-2.2.3-bin.tar.gz

将安装包上传到hadoop100服务器/opt/download路径下,并进行解压:

$ cd /opt/download
$ tar -zxvf hbase-2.2.3-bin.tar.gz -C /opt/pkg/

配置HBase

修改HBase配置文件hbase-env.sh

$ cd /opt/pkg/hbase-2.2.3/conf
$ vim hbase-env.sh

修改如下两项内容,值如下

export JAVA_HOME=/opt/pkg/java
export HBASE_MANAGES_ZK=false  

修改文件hbase-site.xml

$ vim hbase-site.xml

内容如下

<configuration>
    <!-- 指定hbase在HDFS上存储的路径 -->
    <property>
         <name>hbase.rootdir</name>
        <value>hdfs://hadoop100:8020/hbase</value>
    </property>

    <!-- 指定hbase是否分布式运行 -->
    <property>
        <name>hbase.cluster.distributed</name>
        <value>true</value>
    </property>

    <!-- 指定zookeeper的地址,多个用“,”分割 -->
    <property>
        <name>hbase.zookeeper.quorum</name>
         <value>hadoop100:2181</value>
    </property>

    <!--指定hbase管理页面-->
    <property>
       <name>hbase.master.info.port</name>
       <value>16010</value>
    </property>

    <!-- 在分布式的情况下一定要设置,不然容易出现Hmaster起不来的情况 -->
    <property>
        <name>hbase.unsafe.stream.capability.enforce</name>
        <value>false</value>
    </property>
</configuration>

修改regionservers配置文件,指定HBase的从节点主机名:

$ vim regionservers
hadoop100

添加HBase环境变量

export HBASE_HOME=/opt/pkg/hbase-2.2.3
export PATH=$PATH:$HBASE_HOME/bin

重新执行/etc/profile,让环境变量生效

source /etc/profile 

HBase的启动与停止:

启动HBase前需要提前启动HDFS及ZooKeeper集群:

如果没开启hdfs,请在运行命令:

$ start-dfs.sh

如果没开启zookeeper,请运行命令:

$ zkServer.sh start conf/zoo.cfg

执行以下命令启动HBase集群

$ start-hbase.sh

启动完后,jps查看HBase相关进程是否都正常运行:

$ jps
7601 QuorumPeerMain
1670 NameNode
1975 SecondaryNameNode
6695 HMaster
6841 HRegionServer
1787 DataNode
2491 JobHistoryServer
2236 ResourceManager
6958 Jps
2351 NodeManager

使用HBase提供的Shell客户端进行访问:

$ hbase shell
HBase Shell
Use "help" to get list of supported commands.
Use "exit" to quit this interactive shell.
For Reference, please visit: http://hbase.apache.org/2.0/book.html#shell
Version 2.2.3, r6a830d87542b766bd3dc4cfdee28655f62de3974, 2020年 01月 10日 星期五 18:27:51 CST
Took 0.0025 seconds
hbase(main):001:0> status

访问HBase的WEB UI界面:

浏览器页面访问 http://hadoop100:16010

image-20220210023039865

停止HBase相关进程的命令:

stop-hbase.sh

安装MySQL

下载rpm-bundle包

$ wget http://mirrors.163.com/mysql/Downloads/MySQL-5.7/mysql-5.7.33-1.el7.x86_64.rpm-bundle.tar
$ tar xvf mysql-5.7.33-1.el7.x86_64.rpm-bundle.tar
$ mkdir mysql-jars
$ mv mysql-comm*.rpm mysql-jars/
$ cd mysql-jars/
$ ls
mysql-community-client-5.7.33-1.el7.x86_64.rpm
mysql-community-common-5.7.33-1.el7.x86_64.rpm
mysql-community-libs-5.7.33-1.el7.x86_64.rpm
mysql-community-libs-compat-5.7.33-1.el7.x86_64.rpm
mysql-community-server-5.7.33-1.el7.x86_64.rpm

依次手动安装

sudo rpm -ivh mysql-community-common-5.7.33-1.el7.x86_64.rpm
sudo rpm -ivh mysql-community-libs-5.7.33-1.el7.x86_64.rpm
sudo rpm -ivh mysql-community-libs-compat-5.7.33-1.el7.x86_64.rpm
sudo rpm -ivh mysql-community-client-5.7.33-1.el7.x86_64.rpm
sudo rpm -ivh mysql-community-server-5.7.33-1.el7.x86_64.rpm

启动MySQL服务

启动MySQL服务

$ systemctl start mysqld.service

$ service mysqld start

配置开机启动

$ systemctl enable mysqld.service

查看运行状态

$ systemctl status mysqld.service

mysql   2574   1 1 23:49 ?    00:00:00 /usr/sbin/mysqld --daemonize --pid-file=/var/run/mysqld/mysqld.pid

查看到进程信息

$ sudo netstat -anpl | grep mysql
tcp6    0   0 :::3306         :::*          LISTEN   946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45824  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45818  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45838  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45808  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45816  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45820  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45826  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45822  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45814  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45846  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45828  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45832  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45842  ESTABLISHED 946/mysqld
tcp6    0    0 192.168.186.100:3306  192.168.186.100:45830  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45844  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45812  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45836  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306   192.168.186.100:45834  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45840  ESTABLISHED 946/mysqld
tcp6    0   0 192.168.186.100:3306  192.168.186.100:45810  ESTABLISHED 946/mysqld
unix 2   [ ACC ]   STREAM   LISTENING   20523  946/mysqld      /var/lib/mysql/mysql.sock

查看端口信息:

可以看出mysql server的进程mysqld所使用的默认端口即3306

$ sudo netstat -anpl | grep tcp
tcp    0   0 0.0.0.0:22       0.0.0.0:*        LISTEN   1412/sshd     
tcp    0   52 192.168.186.103:22   192.168.186.1:54058   ESTABLISHED 2287/sshd: hadoop 
tcp6    0   0 :::3306         :::*          LISTEN   2574/mysqld    
tcp6    0   0 192.168.186.103:3888  :::*          LISTEN   2060/java     
tcp6    0   0 :::22          :::*          LISTEN   1412/sshd     
tcp6    0   0 :::37791        :::*          LISTEN   2060/java     
tcp6    0    0 :::2181         :::*          LISTEN   2060/java 

修改密码

找到临时密码:

第一次登陆mysql需要root的临时密码,这个密码是安装时随机生成在MySQL的服务器日志中的:

$ grep "temporary password" /var/log/mysqld.log
2020-02-26T17:05:45.104999Z 1 [Note] A temporary password is generated for root@localhost: bl/!6qaU.wuX
$ mysql -u roop -p 

这时输入临时密码即可登录MySQL客户端。

修改密码安全策略:

第一次登陆MySQL客户端终端后系统会很快提示你修改掉默认密码

mysql> show databases;
ERROR 1820 (HY000): You must reset your password using ALTER USER statement before executing this statement.

mysql> alter user 'root'@'localhost' identified by 'niit1234';
ERROR 1819 (HY000): Your password does not satisfy the current policy requirements

基于默认密码安全策略,所设置的密码必须要包含大小写字母、数字和字符。如果不考虑安全问题,可以修改策略:

mysql> set global validate_password_policy=0;
mysql> set global validate_password_length=1;

设置新root密码并开启远程权限:

Mysql客户端远程访问因为安全原因默认是关闭的。我们需要将root的访问权限扩大都允许从任意ip访问:

mysql>GRANT ALL PRIVILEGES ON *.* TO 'root'@'%'IDENTIFIED BY 'niit1234' WITH GRANT OPTION;

现在虽然修改了远程访问权限,但是还没有生效,因此我们需要刷新权限:

mysql>FLUSH PRIVILEGES;

修改默认字符集

修改MySQL配置文件:

vi /etc/my.cnf

在文件后面追加:

# 默认服务器内部操作字符集
character-set-server=utf8mb4
# 默认服务器内部操作字符集校对规则
collation-server=utf8mb4_general_ci
# 默认的存储引擎
default-storage-engine=InnoDB
# 初始化连接时设置以下字符集:
# character_set_client
# character_set_results
# character_set_connection
init_connect='set names utf8mb4'

[client]
default-character-set=utf8mb4

[mysql]
default-character-set=utf8mb4

重启MySQL服务

$ sudo service mysqld restart

安装Hive

下载安装包

从官网下载hive安装包:apache-hive-3.1.2-bin.tar.gz

规划安装目录:/opt/pkg/hive

上传安装包到hadoop100服务器中

解压到安装路径

解压安装包到指定的规划目录/opt/pkg/

$ cd /export/softwares/
$ tar -xzvf apache-hive-3.1.2-bin.tar.gz -C /opt/pkg/

修改配置文件

进入hive安装目录:

$ cd /opt/pkg

重新命名hive目录:

$ mv apache-hive-3.1.2-bin/hive

修改/opt/pkg//conf目录下的hive-site.xml,默认没有该文件, 需要手动创建:

$ cd /opt/pkg/hive/conf/
$ vim hive-site.xml

进入编辑模式, 文件内容如下:

<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
   <property>
       <name>javax.jdo.option.ConnectionURL</name>
       <value>jdbc:mysql://hadoop100:3306/metastore?useSSL=false</value>
   </property>

   <property>
       <name>javax.jdo.option.ConnectionDriverName</name>
       <value>com.mysql.jdbc.Driver</value>
   </property>

   <property>
       <name>javax.jdo.option.ConnectionUserName</name>
       <value>root</value>
   </property>

   <property>
       <name>hive.metastore.warehouse.dir</name>
       <value>/user/hive/warehouse</value>
   </property>

   <property>
       <name>javax.jdo.option.ConnectionPassword</name>
       <value>niit1234</value>
   </property>

   <property>
       <name>hive.metastore.schema.verification</name>
       <value>false</value>
   </property>

   <property>
       <name>hive.metastore.event.db.notification.api.auth</name>
       <value>false</value>
   </property>

    <property>
       <name>hive.cli.print.current.db</name>
       <value>true</value>
   </property>

    <property>
       <name>hive.cli.print.header</name>
       <value>true</value>
   </property>

   <property>
       <name>hive.server2.thrift.bind.host</name>
       <value>hadoop100</value>
   </property>

   <property>
       <name>hive.server2.thrift.port</name>
       <value>10000</value>
   </property>
</configuration>

创建hive日志存储目录

$ mkdir /opt/pkg/hive/logs/

重命名日志配置文件模板为hive-log4j.properties

$ pwd
/opt/pkg/hive/conf

$ mv hive-log4j2.properties.template hive-log4j2.properties
$ vim hive-log4j2.properties # 修改文件

修改此文件的hive.log.dir属性的值:

#更改以下内容,设置我们的hive的日志文件存放的路径,便于排查问题

hive.log.dir=/opt/pkg/hive/logs/

拷贝mysql驱动包

由于运行hive时,需要向mysql数据库中读写元数据,所以需要将mysql的驱动包上传到hive的lib目录下。

上传mysql驱动包,如mysql-connector-java-5.1.38.jar/opt/download/目录中:

$ cp mysql-connector-java-5.1.38.jar /opt/pkg/hive/lib/

解决日志Jar包冲突

# 进入lib目录
$ cd /opt/pkg/hive/lib/

# 重新命名 或者直接删除
$ mv log4j-slf4j-impl-2.10.0.jar log4j-slf4j-impl-2.10.0.jar.bak

配置Hive环境变量

# HIVE
export HIVE_HOME=/opt/pkg/hive
export PATH=$PATH:$HIVE_HOME/bin

# HCATLOG
export HCAT_HOME=$HIVE_HOME/hcatalog
export PATH=$PATH:$HCAT_HOME/bin

之后别忘记source使环境变量配置文件的修改生效

初始化元数据库

开启MySQL客户端连接MySQL服务, 用户名root, 密码niit1234,创建hive元数据库, 数据库名称需要和hive-site.xml中配置的一致:

$ mysql -uroot -pniit1234
create database metastore;
show databases;

退出mysql:

Exit;

初始化元数据库:

$ schematool -initSchema -dbType mysql -verbose

看到schemaTool completed 表示初始化成功

启动Hive服务

前提:Hadoop集群、MySQL服务均已启动,执行命令:

nohup hive --service metastore >/tmp/metastore.log 2>&1 &
nohup hive --service hiveserver2 >/tmp/hiveServer2.log 2>&1 &

验证Hive安装是否成功

在hadoop100上任意目录启动hive的命令行客户端beeline:

$ beeline
Beeline version 3.1.2 by Apache Hive
beeline> !connect jdbc:hive2://localhost:10000
Connecting to jdbc:hive2://localhost:10000
Enter username for jdbc:hive2://localhost:10000: hadoop
Enter password for jdbc:hive2://localhost:10000: ******
Connected to: Apache Hive (version 3.1.2)
Driver: Hive JDBC (version 3.1.2)
Transaction isolation: TRANSACTION_REPEATABLE_READ
0: jdbc:hive2://localhost:10000> show databases;
+----------------+
| database_name |
+----------------+
| default    |
+----------------+
2 rows selected (1.93 seconds)

认证时密码直接回车即可, 如果能看到以上信息说明hive安装成功:

退出客户端:

0: jdbc:hive2://localhost:10000> !quit
Closing: 0: jdbc:hive2://localhost:10000 quit;

使用beeline连接失败的解决办法

如果hiveserver2已经正常运行在本机的10000端口上, 但使用beeline连接hiveserver2报错, WARN jdbc.HiveConnection: Failed to connect to localhost:10000
主要的原因可能是hadoop引入了用户代理机制, 不允许上层系统直接使用实际用户.

解决办法: 在hadoop的核心配置文件core-site.xml中添加:

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

其中hadoop.proxyuser.XXX的XXX就是代理用户名,我这里设置成hadoop
然后重启hadoop,否则以上修改不生效。

$ stop-all.sh

$ start-dfs.sh
$ start-yarn.sh

再次执行下面命令就可以正常连接欸蓝:

beeline -u jdbc:hive2://localhost:10000

安装Sqoop

Sqoop下载和安装

下载地址:http://archive.apache.org/dist/sqoop/1.4.7/

wget http://archive.apache.org/dist/sqoop/1.4.7/sqoop-1.4.7.bin__hadoop-2.6.0.tar.gz

解压sqoop-1.4.7.bin__hadoop-2.6.0.tar.gz到指定目录下:

tar -zxvf sqoop-1.4.7.bin__hadoop-2.6.0.tar.gz -C /opt/pkg/

修改Sqoop安装目录名称:

$ cd /opt/pkg/
$ mv sqoop-1.4.7.bin__hadoop-2.6.0/ sqoop

配置Sqoop环境变量

$ vi ~/.bash_profile

#sqoop
export SQOOP_HOME=/opt/pkg/sqoop
export PATH=$PATH:$SQOOP_HOME/bin

把Sqoop所依赖的相关环境变量都配置上,修改sqoop-env.sh

$ mv conf/sqoop-env-template.sh conf/sqoop-env.sh
$ vi conf/sqoop-env.sh 

# Set Hadoop-specific environment variables here.

#Set path to where bin/hadoop is available
export HADOOP_COMMON_HOME=/opt/pkg/hadoop

#Set path to where hadoop-*-core.jar is available
export HADOOP_MAPRED_HOME=/opt/pkg/hadoop

#set the path to where bin/hbase is available
export HBASE_HOME=/opt/pkg/hbase

#Set the path to where bin/hive is available
export HIVE_HOME=/opt/pkg/hive

#Set the path for where zookeper config dir is
export ZOOCFGDIR=/opt/pkg/zookeeper/conf

修改sqoop-site.xml, 具体配置如下文件所示:

<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
  <property>
  <name>sqoop.metastore.client.enable.autoconnect</name>
  <value>true</value>
  <description>If true, Sqoop will connect to a local metastore
   for job management when no other metastore arguments are
   provided.
  </description>
 </property>
 <property>
  <name>sqoop.metastore.client.autoconnect.url</name>
  <value>jdbc:hsqldb:file:/tmp/sqoop-meta/meta.db;shutdown=true</value>
  <description>The connect string to use when connecting to a
   job-management metastore. If unspecified, uses ~/.sqoop/.
   You can specify a different path here.
  </description>
 </property>
 <property>
  <name>sqoop.metastore.client.autoconnect.username</name>
  <value>SA</value>
  <description>The username to bind to the metastore.
  </description>
 </property>
 <property>
  <name>sqoop.metastore.client.autoconnect.password</name>
  <value></value>
  <description>The password to bind to the metastore.
  </description>
 </property>
  <property>
  <name>sqoop.metastore.client.record.password</name>
  <value>true</value>
  <description>If true, allow saved passwords in the metastore.
  </description>
 </property>
 <property>
  <name>sqoop.metastore.server.location</name>
  <value>/tmp/sqoop-metastore/shared.db</value>
  <description>Path to the shared metastore database files.
  If this is not set, it will be placed in ~/.sqoop/.
  </description>
 </property>
 <property>
  <name>sqoop.metastore.server.port</name>
  <value>16000</value>
  <description>Port that this metastore should listen on.
  </description>
 </property>
</configuration>

修改configure-sqoop

$ vi bin/configure-sqoop

如果没有安装accumulo,则将有关ACCUMULO_HOME的判断逻辑注释掉:

$ vi bin/configure-sqoop
 94 #if [ -z "${ACCUMULO_HOME}" ]; then
 95 # if [ -d "/usr/lib/accumulo" ]; then
 96 #  ACCUMULO_HOME=/usr/lib/accumulo
 97 # else
 98 #  ACCUMULO_HOME=${SQOOP_HOME}/../accumulo
 99 # fi
100 #fi

140 #if [ ! -d "${ACCUMULO_HOME}" ]; then
141 # echo "Warning: $ACCUMULO_HOME does not exist! Accumulo imports will fail."
142 # echo 'Please set $ACCUMULO_HOME to the root of your Accumulo installation.'
143 #fi

这样做的目的是避免将来运行时出现类似下面的警告信息:

Warning: /opt/pkg/sqoop/bin/…/…/accumulo does not exist! Accumulo imports will fail.
Please set $ACCUMULO_HOME to the root of your Accumulo installation.

将MySQL的驱动(MySQL5.7对应的驱动版本应为5.x版本)上传到Sqoop安装目录下的lib目录下:

$ cp mysql-connector-java-5.1.44-bin.jar /opt/pkg/sqoop/lib/

$HIVE_HOME/lib/hive-common-3.1.2.jar拷贝或者软链接到$SQOOP_HOME/lib

$ ln -s /opt/pkg/hive/lib/hive-common-3.1.2.jar /opt/pkg/sqoop/lib

如果需要解析json,可下载java-json.jar放到sqoop目录下的lib里。

下载地址:http://www.java2s.com/Code/Jar/j/Downloadjavajsonjar.htm

$ cp java-json.jar /opt/pkg/sqoop/lib/

如果需要avro序列化,可将hadoop里面的avro的jar包拷贝或者软链接到sqoop目录下的lib里。

$ ln -s /opt/pkg/hadoop/share/hadoop/common/lib/avro-1.7.7.jar /opt/pkg/sqoop/lib/

练习

完成以下练习:

练习1:

安装好Sqoop学习环境后使用cd命令进入到sqoop安装目录, 输入以下命令并观察输出:

$ sqoop version
2020-01-23 23:43:20,287 INFO sqoop.Sqoop: Running Sqoop version: 1.4.7
Sqoop 1.4.7
git commit id 2328971411f57f0cb683dfb79d19d4d19d185dd8
Compiled by maugli on Thu Dec 21 15:59:58 STD 2017

如果输出了类似内容, 表明Sqoop的安装初步完成.

Views: 917

Kafka-Storm 实时计算项目开发实战

项目架构

image-20211209020000874

基本要求

  1. 主题相关的WebAPP
    1. 1个首页和若干功能页面
    2. 有生成数据的能力(实际展示的时候,为了有更多丰富的数据, 允许使用脚本生成假数据)
      1. KafkaProducer 直接将需要采集的信息发送到Storm
    3. 分析结果的图表展示
      1. echarts, 或者其他图表库
        1. Jquery的ajax库
        2. json
    4. 其他要求:
      1. 参照Alibaba的Java开发手册的规约
      2. 代码中适当添加注释
  2. 部署到服务器(虚拟机或者购买的服务器)
    1. JavaWeb程序需要部署到 Tomcat
    2. Nginx 实现反向代理
  3. 实时数据分析
    1. kafka -> storm
    2. 持久化到数据库(HBase,MySQL,Redis)
    3. 最好每个团队成员都有一个完整的实时分析流程
    4. 每个团队最少有两种实时分析, 不同类型(不要过于简单,分析的内容要能体现实时性)

开发流程建议

  1. 确定需求和分工

  2. 学习使用版本控制工具(Gitee 码云 / Github)

  3. 统一消息格式,建议使用JSON

  4. 开发时先采用本地的Tomcat和Storm环境测试

  5. 本地环境测试时可以远程连接服务器上的的Kafka和Hbase

    没问题再使用Tomcat和Nginx部署到服务器

  6. 将拓扑上传到Storm集群中运行

  7. 联调所有模块

  8. 完善SRS报告

    记录整个项目的需求, 架构, 开发、部署、测试的详细过程

  9. 每个人准备PPT, 视频等交付资料

  10. 每周使用在线表格记录项目进度情况

WebApp开发

采用前后端分离的模式开发和部署

后端可以采用基于tomcat的普通JavaWeb应用或者SpringBoot项目

消息传递

一般项目可以使用kafka-clients依赖里的KafkaProducer API来发送消息到Kafka, 这样便于快速调试.

对于SpringBoot项目,可以使用spring-kafka依赖.

kafka相关配置

spring:
  kafka:
    bootstrap-servers: hadoop000:9092,hadoop000:9093,hadoop000:9094
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      client-id: app-pro-cli
      acks: 1
      retries: 3
    consumer:
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      client-id: app-pro-cli
      group-id: g1

Kafka的初始化配置

package cn.delucia.project.conf;

import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class KafkaInitConf {
    @Value("${app.kafka-topic}")
    private String kafkaTopic;
    @Value("${app.topic-partitions}")
    private Integer topicPartitions;
    // 创建一个Topic并设置分区数和副本数
    @Bean
    public NewTopic initialTopic() {
        return new NewTopic(kafkaTopic, 1, (short) 1);
    }
    // 如果要修改分区数,只需修改配置值重启项目即可
    // 修改分区数并不会导致数据的丢失,但是分区数只能增大不能减小
    @Bean
    public NewTopic updateTopic() {
        return new NewTopic(kafkaTopic, topicPartitions, (short) 1);
    }
}

在控制器类中注入kafkaTemplate并调用send方法发送消息, 例如:

@Slf4j
@RestController
public class GreetingController {

    @Autowired
    private KafkaTemplate<Object, String> kafkaTemplate;

    public Greeting greeting(@RequestParam(value = "name", defaultValue = "农场主") String name) {

        Greeting greeting = new Greeting(counter.incrementAndGet(), name);

        try {
            String s = new ObjectMapper().writeValueAsString(greeting);
            // 带回调的生产者
            kafkaTemplate.send(kafkaTopic, "Greeting:" + s).addCallback(success -> {
                // 消息发送到的topic
                String topic = Objects.requireNonNull(success).getRecordMetadata().topic();
                // 消息发送到的分区
                int partition = success.getRecordMetadata().partition();
                // 消息在分区内的offset
                long offset = success.getRecordMetadata().offset();
                log.info("发送消息成功: {}-{}-{}, Greeting:{}", topic, partition, offset, s);
            }, failure -> {
                log.info("发送消息失败: {}", failure.getMessage());
            });
        } catch (JsonProcessingException e) {
            log.info("解析Json格式失败: {}", e.getMessage());
        }
        ...
    }
}

安装Tomcat

准备:安装Java JDK1.8

这里安装openjdk是因为比较简单, 工作场合一定要使用oracle提供的jdk1.8

sudo yum -y install java-1.8.0-openjdk*

这样安装的好处就是环境变量都配好了

可以直接查看版本 java -version

Tomcat安装

下载页面: https://tomcat.apache.org/download-90.cgi

文档:https://tomcat.apache.org/tomcat-9.0-doc/index.html

下载链接

wget https://mirrors.tuna.tsinghua.edu.cn/apache/tomcat/tomcat-9/v9.0.41/bin/apache-tomcat-9.0.41.tar.gz

解压到后改名tomcat9

tar -zxvf apache-tomcat-9.0.41.tar.gz -C ~/app
cd ~/app
mv apache-tomcat-9.0.41/ tomcat9

默认tomcat端口是8080, 为了避免冲突, 这里修改为18080

 $ vi ~/app/tomcat9/conf/server.xml

 69     <Connector port="18080" protocol="HTTP/1.1"
 70                connectionTimeout="20000"
 71                redirectPort="8443" />

启动: 进入 /bin目录下 运行startup.sh脚本文件

$ ./startup.sh 
Using CATALINA_BASE:   /home/hadoop/app/tomcat9
Using CATALINA_HOME:   /home/hadoop/app/tomcat9
Using CATALINA_TMPDIR: /home/hadoop/app/tomcat9/temp
Using JRE_HOME:        /home/hadoop/app/jdk1.8.0_211
Using CLASSPATH:       /home/hadoop/app/tomcat9/bin/bootstrap.jar:/home/hadoop/app/tomcat9/bin/tomcat-juli.jar
Tomcat started.

检查tomcat进程信息

$ bin]$ ps -ef | grep tomcat
hadoop    15142      1  1 14:41 pts/0    00:00:05 /home/hadoop/app/jdk1.8.0_211/bin/java -Djava.util.logging.config.file=/home/hadoop/app/tomcat9/conf/logging.properties -Djava.util.logging.manager=org.apache.juli.ClassLoaderLogManager -Djdk.tls.ephemeralDHKeySize=2048 -Djava.protocol.handler.pkgs=org.apache.catalina.webresources -Dorg.apache.catalina.security.SecurityListener.UMASK=0027 -Dignore.endorsed.dirs= -classpath /home/hadoop/app/tomcat9/bin/bootstrap.jar:/home/hadoop/app/tomcat9/bin/tomcat-juli.jar -Dcatalina.base=/home/hadoop/app/tomcat9 -Dcatalina.home=/home/hadoop/app/tomcat9 -Djava.io.tmpdir=/home/hadoop/app/tomcat9/temp org.apache.catalina.startup.Bootstrap start
hadoop    15296  15064  0 14:47 pts/0    00:00:00 grep --color=auto tomcat

检查对应的监听端口信息

netstat -anpt | grep 15142

网页访问虚拟机主机名或IP地址:18080

为了方便可以为tomcat的安装目录配置环境变量CATALINA_HOME,添加到PATH中

为tomcat添加用户和角色

修改conf/tomcat-users.xml, 添加如下内容

<!--  
  <role rolename="tomcat"/>
  <role rolename="role1"/>
  <user username="tomcat" password="tomcat" roles="tomcat"/>
  <user username="both" password="tomcat" roles="tomcat,role1"/>
  <user username="role1" password="tomcat" roles="role1"/>
-->
  <role rolename="manager-gui"/>
  <role rolename="manager-script" />
  <user username="admin" password="123123" roles="manager-gui,manager-script" />
</tomcat-users>

修改webapps/manager/META-INF目录下的context.xml,在allow行的末尾加上|\d+.\d+.\d+.\d+表示允许所有主机访问。

<Context antiResourceLocking="false" privileged="true" >
  <Valve className="org.apache.catalina.valves.RemoteAddrValve"
         allow="127\.\d+\.\d+\.\d+|::1|0:0:0:0:0:0:0:1|\d+\.\d+\.\d+\.\d+" />
  <Manager sessionAttributeValueClassNameFilter="java\.lang\.(?:Boolean|Integer|Long|Number|String)|org\.apache\.catalina\.filters\.CsrfPreventionFilter\$LruCache(?:\$1)?|java\.util\.(?:Linked)?HashMap"/>
</Context>

重启tomcat生效

部署项目

一般的javaee项目可以直接build出一个war包进行上传服务器。

SpringBoot默认是打成jar包,如果需要打成war包, 需要修改pom.xml文件:

    <groupId>com.niit</groupId>
    <artifactId>demo</artifactId>
    <version>0.0.1-SNAPSHOT</version>

    <!-- 这里打成war包 若打jar,需将war改为jar -->
    <packaging>war</packaging>

    <name>demo</name>
    <description>Demo project for Spring Boot</description>

然后使用mvn:package构建war包即可,然后上传到服务器的tomcat目录下的webapp文件夹之内。

tomcat 容器的运行机制👇

tomcat默认会加载tomcat目录下的webapp文件夹之内的文件,如

其中ROOT目录下为Tomcat的欢迎页

http://hadoop000/index.jsp

http://hadoop000/tomcat.gif

examples目录是一些官方示例

http://hadoop000/examples/

tomcat也会默认会加载tomcat目录下的webapp文件夹中下面的war包,并自动解压在webapp下面。

默认的访问方式就是 http://域名:端口号/war包名, 端口号默认是8080

启动tomcat, 浏览器访问 http://hadoop000:18080/demo/

部署SpringBoot项目

项目的服务器设置:

server:
  port: 18080
  servlet:
    context-path: "/project"
#debug: on

使用maven的springboot打包插件按照jar包的方式打包

   <groupId>com.niit</groupId>
    <artifactId>demo</artifactId>
    <version>0.0.1-SNAPSHOT</version>

    <!-- 这里打成war包 若打jar,需将war改为jar -->
    <packaging>jar</packaging>

    <name>demo</name>
    <description>Demo project for Spring Boot</description>

然后使用mvn:package打包即可,将打包出来的jar文件重命名然后上传到服务器

在服务器启动项目(需要Java8以上环境)

[hadoop@hadoop000 webapps]$ java -jar demo-jar.jar 

  .   ____          _            __ _ _
 /\\ / ___'_ __ _ _(_)_ __  __ _ \ \ \ \
( ( )\___ | '_ | '_| | '_ \/ _` | \ \ \ \
 \\/  ___)| |_)| | | | | || (_| |  ) ) ) )
  '  |____| .__|_| |_|_| |_\__, | / / / /
 =========|_|==============|___/=/_/_/_/
 :: Spring Boot ::        (v2.3.0.RELEASE)

如果后台运行可以(推荐)

$ nohup java -jar demo-jar.jar 1> demo.log 2>&1 &

浏览器访问 http://hadoop000:18080/project/greeting

配置Nginx反向代理

但是这样带8080端口的访问并不是很好,因此一般都使用nginx反向代理.

通过反向代理将对nginx的80端口的访问请求转发到tomcat的8080端口

首先在物理机hosts文件创建虚拟机主机名和ip的映射

192.168.186.100 hadoop000

Nginx安装

安装前确认是否已经安装过

sudo yum search nginx

Nginx文档http://nginx.org/en/docs/

Installation on Linux, nginx packages from nginx.org can be used.

Installation instructions
RHEL/CentOS
Debian
Ubuntu
SLES
Alpine

以 CentOS为例进行安装

Install the prerequisites:

sudo yum install yum-utils

To set up the yum repository, create the file named /etc/yum.repos.d/nginx.repo with the following contents:

sudo vi /etc/yum.repos.d/nginx.repo

内容如下:

[nginx-stable]
name=nginx stable repo
baseurl=http://nginx.org/packages/centos/$releasever/$basearch/
gpgcheck=1
enabled=1
gpgkey=https://nginx.org/keys/nginx_signing.key
module_hotfixes=true

[nginx-mainline]
name=nginx mainline repo
baseurl=http://nginx.org/packages/mainline/centos/$releasever/$basearch/
gpgcheck=1
enabled=0
gpgkey=https://nginx.org/keys/nginx_signing.key
module_hotfixes=true

By default, the repository for stable nginx packages is used. If you would like to use mainline nginx packages, run the following command:

sudo yum-config-manager --enable nginx-mainline

To install nginx, run the following command:

sudo yum install -y nginx

When prompted to accept the GPG key, verify that the fingerprint matches 573B FD6B 3D8F BC64 1079 A6AB ABF5 BD82 7BD9 BF62, and if so, accept it.

Nginx配置

创建java.conf ,进行最简配置

$ cd /etc/nginx/conf.d/
$ sudo cp default.conf java.conf
$ vi java.conf

URL记住不要忘了加http://前缀

server {
    listen  80;
    server_name hadoop000;
    location / {
        proxy_pass http://127.0.0.1:18080;
    }
}

常用命令

解释 命令
安装服务 yum install nginx
启动服务 service nginx start
停止服务 service nginx stop
重载服务 service nginx reload

配置完成后启动服务

sudo service nginx start

如果服务器已经启动,当配置发生变化可以直接使用重载服务来更新配置,运维常用,因为不需要停止服务就可以重载新的配置。

如果启动失败可以查看错误日志

sudo vi /var/log/nginx/error.log 

经过反向代理配置之后,使用浏览器访问 http://hadoop000就相当于访问虚拟机hadoop000的本地服务http://127.0.0.1:18080的效果

如果发现502错误:

2020/11/18 16:12:39 [crit] 16376#16376: *1 connect() to 127.0.0.1:8080 failed (13: Permission denied) while connecting to upstream, client: 192.168.186.1, server: hadoop000, request: "GET /favicon.ico HTTP/1.1", upstream: "http://127.0.0.1:8080/favicon.ico", host: "hadoop000", referrer: "http://hadoop000/"

此时需要考虑把linux操作系统默认的强制访问安全限制设置为禁用。

关闭SElinux即可

  1. 临时关闭 SElinux

    sudo setenforce 0
  2. 永久关闭 SElinux

    sudo vim /etc/selinux/config
    SELINUX=disabled

修改之后就可以正常访问了

配置开机启动

[hadoop@hadoop000 download]$ chkconfig nginx
注意:正在将请求转发到“systemctl is-enabled nginx.service”。
disabled

[hadoop@hadoop000 download]$ chkconfig nginx on
注意:正在将请求转发到“systemctl enable nginx.service”。
==== AUTHENTICATING FOR org.freedesktop.systemd1.manage-unit-files ===
Authentication is required to manage system service or unit files.
Authenticating as: root
Password: 
==== AUTHENTICATION COMPLETE ===
Created symlink from /etc/systemd/system/multi-user.target.wants/nginx.service to /usr/lib/systemd/system/nginx.service.
==== AUTHENTICATING FOR org.freedesktop.systemd1.reload-daemon ===
Authentication is required to reload the systemd state.
Authenticating as: root
Password: 
==== AUTHENTICATION COMPLETE ===

[hadoop@hadoop000 download]$ chkconfig nginx
注意:正在将请求转发到“systemctl is-enabled nginx.service”。
enabled

前后端分离

nginx是一个高性能服务器,除了配置反向代理之外,也非常适合部署静态资源并支持高并发请求, 并且也提供负载均衡的功能。

首先删除默认的配置

# sudo vim /etc/nginx/conf.d/default.conf

为了将反向代理请求和静态资源的请求分开,修改我们之前的配置如下:

upstream webapp.server {
    server hadoop000:18080;
}

server {
        listen  80;
        server_name localhost hadoop000;
        root /data/www/;

        # 静态资源
        location / {
            index index.html;
            access_log /var/log/nginx/java-hadoop.log main;
        }

        # 反向代理到本地JavaWeb的后台服务
        location ^~ /project/ {
            proxy_pass http://webapp.server/project/;
        }
}

一些说明如下:

  1. listen 80 是http协议的默认端口,:80 可以省略
  2. server_name 表示请求路径中的服务器主机名,可配置多个
  3. /data/www 为站点根目录,需手动创建,权限一般为755
  4. access_log 对应的路径是应用的服务器日志,文件夹不存在则需手动创建
  5. location 的匹配规则,优先匹配 /demo/,其次 /

一般来说前后端分离部署应该是部署在不同的服务器上的,这里放在一台服务器上只是为了演示方便。

前端部署

然后将项目的静态页面放到/data/www下即可

后端部署

SpringBoot项目移除静态资源,单独部署在tomcat上,并使用ngixn实现反向代理

实时数据采集

主要技术:Kafka

image-20211209020033601

实时数据流计算

主要技术:Kafka-clients,storm-hbase|storm-redis|storm-mysql

以热力图项目为例, 拓扑代码如下

package com.niit.project;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.StormSubmitter;
import org.apache.storm.hbase.bolt.HBaseBolt;
import org.apache.storm.hbase.bolt.mapper.SimpleHBaseMapper;
import org.apache.storm.kafka.spout.ByTopicRecordTranslator;
import org.apache.storm.kafka.spout.KafkaSpout;
import org.apache.storm.kafka.spout.KafkaSpoutConfig;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

import java.util.HashMap;
import java.util.Map;

public class KafkaStormProjectApp {

    public static String topologyName = "project-topo";
    public static final String KAFKA_BROKER = "hadoop000:9092";
    public static final String INPUT_TOPIC = "storm-project";

    private static class SplitBolt extends BaseRichBolt {
        private OutputCollector outputCollector;

        @Override
        public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) {
            this.outputCollector = collector;
        }

        @Override
        public void execute(Tuple input) {
            String id = input.getStringByField("line");
            String[] split = id.split(",");

            try {
                double lng = Double.parseDouble(split[0]);
                double lat = Double.parseDouble(split[1]);
                this.outputCollector.emit(new Values(id, lng, lat));
            } catch (NumberFormatException e) {
                System.err.println(e.getMessage());
            }
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("id", "lng", "lat"));
        }
    }

    /**
     * 计数利用HBase的CountColumn特性
     */
    private static class CountBolt extends BaseRichBolt {

        private OutputCollector collector;
        private final HashMap<String, Long> counts = null;

        @Override
        public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) {
            this.collector = collector;
        }

        @Override
        public void execute(Tuple input) {
            String id = input.getString(0);
            double lng = input.getDouble(1);
            double lat = input.getDouble(2);

            this.collector.emit(new Values(id, lng, lat, 1L));
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("id", "lng", "lat", "count"));
        }
    }

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

        final TopologyBuilder builder = new TopologyBuilder();

        // storm conf
        Config conf = new Config();
        conf.setNumAckers(0);
        conf.setDebug(true);

        // kafka bolt
        ByTopicRecordTranslator<String, String> translator =
                new ByTopicRecordTranslator<>((r) -> new Values(r.value()), new Fields("line"));
        translator.forTopic(INPUT_TOPIC, (r) -> new Values(r.value()), new Fields("line"));

        KafkaSpoutConfig<String, String> kafkaSpoutConfig = KafkaSpoutConfig
                // bootstrapServers 以及topic
                .builder(KAFKA_BROKER, INPUT_TOPIC)
                // 设置group.id
                .setProp(ConsumerConfig.GROUP_ID_CONFIG, "location")
                // ensure at-least-once processing
                .setProp(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
                // 设置开始消费的起始位置
                // 设置提交消费边界的时长间隔
                .setOffsetCommitPeriodMs(10_000)
                //Translator
                .setRecordTranslator(translator)
                .build();

        KafkaSpout<String, String> kafkaSpout = new KafkaSpout<>(kafkaSpoutConfig);

        // hbase bolt
        Map<String, Object> hbConf = new HashMap<>();
        hbConf.put("hbase.rootdir", "hdfs://hadoop000:9000/hbase");
        hbConf.put("hbase.zookeeper.quorum", "hadoop000:2181");
        conf.put("hbase.conf", hbConf);
        conf.setNumWorkers(2);  // 设置为1个topology创建2个worker进程

        SimpleHBaseMapper mapper = new SimpleHBaseMapper()
                .withRowKeyField("id")
                .withColumnFields(new Fields("lng","lat"))
                .withCounterFields(new Fields("count"))
                .withColumnFamily("cf");

        HBaseBolt hbaseBolt = new HBaseBolt("project", mapper).withConfigKey("hbase.conf");

        // build topology
        builder.setSpout("kafka_spout", kafkaSpout);
        builder.setBolt("split-bolt", new SplitBolt(), 2)
                .setNumTasks(4)
                .shuffleGrouping("kafka_spout");
        builder.setBolt("count-bolt", new CountBolt())
                .fieldsGrouping("split-bolt", new Fields("id"));
        builder.setBolt("hbase-bolt", hbaseBolt).globalGrouping("count-bolt");

        if (args != null && args.length > 0) {
            topologyName = args[0];
            StormSubmitter.submitTopology(topologyName, conf, builder.createTopology());
        } else {
            LocalCluster localCluster = new LocalCluster();
            localCluster.submitTopology(topologyName, conf, builder.createTopology());
        }
    }

}

其中聚合操作是利用了HBase的CounterColumn特性

这里没有使用时间窗口Bolt来体现实时,而是利用了HBase表的TTL属性,TTL可以在创建表的时候指定,TTL设置为60即表示表中记录的存活时间为1分钟:

create "project",{NAME => 'cf', MIN_VERSIONS => '0',TTL => '60'}

也可以disable表之后使用alter语句对已有表进行修改。

由于插入数据的时候经纬度是采用了Double类型,而计数列采用了Long类型,但是HBase只有使用字节数组这样一种方式进行存储,所以需要程序员自己控制数据类型的转换。在查看表中记录的使用需要这样进行数据类型的转换:

scan 'project', {COLUMNS => ['cf:lng:toDouble','cf:lat:toDouble','cf:count:toLong']}

将上传到Storm集群

[kafka-storm-project]$ storm jar project-topology.jar com.niit.demo.KafkaStormProjectTopology project-topo

数据可视化

主要技术:百度echarts图表,异步请求图表渲染(Ajax & JSON)

这个热力图项目的前端比较简单,只有一个index.html文件

<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <title>百度地图</title>
    <style>
        #main {
            position: absolute;
            top: 0;
            left: 0;
            right: 0;
            bottom: 0;
            width: 100%;
            height: 100%;
        }
    </style>
</head>
<body>

<!-- 为ECharts准备一个具备大小(宽高)的Dom -->
<div id="main"></div>

<!-- jquery 1.11.3 -->
<script type="text/javascript" src="https://cdn.jsdelivr.net/npm/jquery@1.11.3/dist/jquery.min.js"></script>
<!-- bootstrap 3.3.7-->
<script type="text/javascript" src="https://cdn.jsdelivr.net/npm/bootstrap@3.3.7/dist/js/bootstrap.min.js"></script>
<!-- echarts 插件  -->
<script type="text/javascript" src="https://cdn.jsdelivr.net/npm/echarts/dist/echarts.min.js"></script>
<!-- echarts 百度地图插件 -->
<script type="text/javascript" src="https://cdn.jsdelivr.net/npm/echarts/dist/extension/bmap.js"></script>
<script type="text/javascript"
        src="https://api.map.baidu.com/api?v=2.0&ak=KOmVjPVUAey1G2E8zNhPiuQ6QiEmAwZu&__ec_v__=20190126"></script>
<script>

    var points = [];
    var myChart = echarts.init(document.getElementById('main'));

    myChart.setOption(option = {
        animation: false,
        bmap: {
            center: [110.337731, 20.064295],  // 海南大学
            zoom: 18, // 地图缩放等级
            roam: true
        },
        visualMap: {
            show: false,
            top: 'top',
            min: 0,
            max: 5,
            seriesIndex: 0,
            calculable: true,
            inRange: {
                color: ['blue', 'blue', 'green', 'yellow', 'red']
            }
        },
        series: [{
            type: 'heatmap',
            coordinateSystem: 'bmap',
            data: points,
            pointSize: 5,
            blurSize: 6
        }]
    });
    // 添加百度地图插件
    var bmap = myChart.getModel().getComponent('bmap').getBMap();
    bmap.addControl(new BMap.MapTypeControl());
    // 禁止拖拽和缩放
    bmap.disableDragging();
    bmap.disableScrollWheelZoom();

    bmap.addEventListener("click", function (e) {
        $.post('log', {
            lng: e.point.lng,
            lat: e.point.lat
        });
    });

    // 10秒更新一次地图
    window.setInterval(function () {
        $.get(
            "points",
            function (data) {
                for (var i = 0; i < data.length; i++) {
                        point = [];
                    for (var j = 0; j < data[i].count; j++) {
                        points.push([data[i].lng, data[i].lat, 1]);
                    }
                }
                option.series.data = points;
                myChart.setOption(option);
            }
        )
    }, 10000);

</script>
</body>
</html>

其中请求demo/points可以到达后端的SpringBoot项目,并且返回HBase的最新数据,对应的接口如下:

src\main\java\com\niit\demo\service\LocationService.java

package com.niit.demo.service;

import com.niit.demo.entity.Point;
import com.niit.demo.utils.HBaseHelper;
import org.apache.hadoop.hbase.Cell;
import org.apache.hadoop.hbase.CellUtil;
import org.apache.hadoop.hbase.client.ResultScanner;
import org.apache.hadoop.hbase.util.Bytes;
import org.springframework.stereotype.Component;

import java.util.ArrayList;
import java.util.List;

@Component
public class LocationService {

    public List<Point> getLocationList() {

        List<Point> list = new ArrayList<>();
        ResultScanner scanner = HBaseHelper.getScanner("project");

        if (scanner != null) {
            scanner.forEach(rowResult -> {

                Point point = new Point();
                for (Cell cell : rowResult.listCells()) {

                    String qualifier = Bytes.toString(CellUtil.cloneQualifier(cell));
                    switch (qualifier) {
                        case "lng":
                            point.setLongitude(Bytes.toDouble(CellUtil.cloneValue(cell)));
                            break;
                        case "lat":
                            point.setLatitude(Bytes.toDouble(CellUtil.cloneValue(cell)));
                            break;
                        case "count":
                            point.setCount(Bytes.toLong(CellUtil.cloneValue(cell)));
                            break;
                        default:
                            break;
                    }
                }
                list.add(point);
            });
        }
        return list;
    }

}

src\main\java\com\niit\demo\IndexController.java

package com.niit.demo;

import com.niit.demo.entity.Point;
import com.niit.demo.service.LocationService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

import java.util.List;

@RestController
@RequestMapping("/")
public class IndexController {
    @Autowired
    public LocationService service;
    Logger logger = LoggerFactory.getLogger(IndexController.class);

    @GetMapping("/points")
    public List<Point> getPoints() {
        return service.getLocationList();
    }

    @PostMapping("/log")
    public void log(double lng, double lat) {
        logger.info("{},{}", lng, lat);
    }
}

常见问题:

Tomcat项目乱码

Tomcat中部署的JSP页面出现中文乱码问题:

  1. 修改tomcat/conf目录下的主配置文件server.xml,添加URIEncoding=“UTF-8”配置项,配置位置如下:

    <Connector port="8080" protocol="HTTP/1.1"
           connectionTimeout="20000"
           redirectPort="8443" URIEncoding="UTF-8" />
    
    <Connector protocol="AJP/1.3"
           address="::1"
           port="8009"
           redirectPort="8443" URIEncoding="UTF-8" />
  2. 修改tomcat/conf/web.xml,在 <servlet>节点中添加如下内容:

    <init-param>
           <param-name>fileEncoding</param-name>
           <param-value>UTF-8</param-value>
    </init-param>

3.重启服务

Tomcat设为开机自启

  1. 进入init.d目录

进入到/etc/init.d目录下,命令是:

cd /etc/init.d
  1. 新建一个名为tomcat的文件
vim tomcat
  1. 为/etc/init.d/tomcat文件添加可执行权限
chmod 755 tomcat
  1. 编辑tomcat文件,添加以下内容
vi tomcat

添加内容为: 注意CATALINA_HOME需要修改正确

#!/bin/bash
# processname: tomcat9
# chkconfig: 2345 86 16
# description: Tomcat9 start|restart|stop.

if [ -f /etc/init.d/functions ]; then
. /etc/init.d/functions
elif [ -f /etc/rc.d/init.d/functions ]; then
. /etc/rc.d/init.d/functions
else
echo -e "/atomcat: unable to locate functions lib. Cannot continue."
exit -1
fi

RETVAL=$?
CATALINA_HOME=/opt/pkg/tomcat9

case "$1" in
start)
if [ -f $CATALINA_HOME/bin/startup.sh ];
then
echo $"Starting Tomcat"
$CATALINA_HOME/bin/startup.sh
fi
;;
stop)
if [ -f $CATALINA_HOME/bin/shutdown.sh ];
then
echo $"Stopping Tomcat"
$CATALINA_HOME/bin/shutdown.sh
fi
;;
*)
echo $"Usage: $0 {start|stop}"
exit 1
;;
esac

exit $RETVAL
  1. 把tomcat这个脚本添加到开机启动项里面
chkconfig --add tomcat
chkconfig tomcat on
  1. 如果想看看是否添加成功
chkconfig --list

netconsole      0:关 1:关 2:关 3:关 4:关 5:关 6:关
network         0:关 1:关 2:开 3:开 4:开 5:开 6:关
tomcat          0:关 1:关 2:开 3:开 4:开 5:开 6:关
  1. 在tomcat/bin下创建一个setenv.sh文件,加入以下环境变量, 并赋予执行权限
[root@hadoop000 tomcat9]# vi bin/setenv.sh

export JAVA_HOME=/opt/pkg/jdk1.8.0_261
export JRE_HOME=/opt/pkg/jdk1.8.0_261/jre
export CATALINA_HOME=/opt/pkg/tomcat9
export CATALINA_BASE=/opt/pkg/tomcat9

[root@hadoop000 tomcat9]# chmod a+x bin/setenv.sh
  1. 查看看是否开机启动

使用命令重启机器,命令是:

reboot
  1. 查看网络状态,执行命令,查看8080端口是否启动
netstat   -lntup
  1. 查看tomcat进程
 ps -ef |grep tomcat
  1. 如果不需要开机启动,从启动脚本删除即可
chkconfig --del tomcat

Views: 418

Index