Spark Streaming 实时流处理引擎

01 什么是Spark Streaming

Spark Streaming是Spark Core API(Spark RDD)的扩展,支持对实时数据流进行可伸缩、高吞吐量、容错处理。数据可以从Kafka、Flume、Kinesis或TCP Socket等多种来源获取,并且可以使用复杂的算法处理数据,这些算法由map()、reduce()、join()和window()等高级函数表示。处理后的数据可以推送到文件系统、数据库等存储系统。事实上,可以将Spark的机器学习和图形处理算法应用于数据流。

image-20220501224946569

Spark Streaming的主要优点如下:

易于使用。Spark Streaming提供了很多高级操作算子,允许以编写批处理作业的方式编写流式作业。它支持Java、Scala和Python语言。

易于与Spark体系整合。通过在Spark Core上运行Spark Streaming,可以在Spark Streaming中使用与Spark RDD相同的代码进行批处理,构建强大的交互应用程序,而不仅仅是数据分析。

02 Spark Streaming工作原理

Spark Streaming接收实时输入的数据流,并将数据流以时间片(秒级)为单位拆分成批次,然后将每个批次交给Spark引擎(Spark Core)进行处理,最终生成以批次组成的结果数据流。

image-20220501225106219

Spark Streaming提供了一种高级抽象,称为DStream(Discretized Stream)。DStream表示一个连续不断的数据流,它可以从Kafka、Flume和Kinesis等数据源的输入数据流创建,也可以通过对其他DStream应用高级函数(例如map()、reduce()、join()和window())进行转换创建。在内部,对输入数据流拆分成的每个批次实际上是一个RDD,一个DStream则由多个RDD组成,相当于一个RDD序列

image-20220501225117153

DStream中的每个RDD都包含来自特定时间间隔的数据。

image-20220501225241853

应用于DStream上的任何操作实际上都是对底层RDD的操作。例如,对一个DStream应用flatMap()算子操作,实际上是对DStream中每个时间段的RDD都执行一次flatMap()算子,生成对应时间段的新RDD,所有的新RDD组成了一个新Dstream。

image-20220501225257562

03 输入DStream和Receiver

输入DStream表示从数据源接收的输入数据流,每个输入DStream(除了文件数据流之外)都与一个Receiver对象相关联,该对象接收来自数据源的数据并将其存储在Spark的内存中进行处理。

如果希望在Spark Streaming应用程序中并行接收多个数据流,可以创建多个输入DStream,同时将创建多个Receiver,接收多个数据流。但需要注意的是,一个Spark Streaming应用程序的Executor是一个长时间运行的任务,它会占用分配给Spark Streaming应用程序的一个CPU内核(占用Spark Streaming应用程序所在节点的一个CPU内核),因此Spark Streaming应用程序需要分配足够的内核(如果在本地运行,则是线程)来处理接收到的数据,并运行Receiver。

在本地运行Spark Streaming应用程序时,不要使用“local”或“local[1]”作为主URL。这两种方式都意味着只有一个线程将用于本地运行任务。如果正在使用基于Receiver的输入DStream(例如Socket、Kafka、Flume等),那么将使用单线程运行Receiver,导致没有多余的线程来处理接收到的数据(Spark Streaming至少需要两个线程,一个线程用于运行Receiver接收数据,一个线程用于处理接收到的数据)。因此,在本地运行时,应该使用“local[n]”作为主URL,其中n>Receiver的数量(若Spark Streaming应用程序只创建了一个DStream,则只有一个Receiver,n的最小值为2)。

每个Spark应用程序都有各自独立的一个或多个Executor进程负责执行任务。将Spark Streaming应用程序发布到集群上运行时,每个Executor进程所分配的CPU内核数量必须大于Receiver的数量,因为1个Receiver独占1个CPU内核,还需要至少1个CPU内核进行数据的处理,这样才能保证至少两个线程同时进行(一个线程用于运行Receiver接收数据,一个线程用于处理接收到的数据)。否则系统将接收数据,但无法进行处理。若Spark Streaming应用程序只创建了一个DStream,则只有一个Receiver,Executor所分配的CPU内核数量的最小值为2。

04 第一个Spark Streaming程序

假设需要监听TCP Socket端口的数据,实时计算接收到的文本数据中的单词数,步骤如下:

  1. 导入相应类

导入Spark Streaming所需的类和StreamingContext中的隐式转换,代码如下:

import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._
  1. 创建StreamingContext

StreamingContext是所有数据流操作的上下文,在进行数据流操作之前需要先创建该对象。例如,创建一个本地StreamingContext对象,使用两个执行线程,批处理间隔为1秒(每隔1秒获取一次数据,生成一个RDD),代码如下:

val conf = new SparkConf()
  .setMaster("local[2]")
  .setAppName("NetworkWordCount")

//按照时间间隔为1秒钟切分数据流
val ssc = new StreamingContext(conf, Seconds(1))
  1. 创建DStream

使用StreamingContext可以创建一个输入DStream,它表示来自TCP源的流数据。例如,从主机名为localhost、端口为9999的TCP源获取数据,代码如下:

val lines = ssc.socketTextStream("localhost", 9999)

上述代码中的lines是一个输入DStream,表示从服务器接收的数据流。lines中的每条记录都是一行文本。

  1. 操作DStream

DStream创建成功后,可以对DStream应用算子操作,生成新的DStream,类似对RDD的操作。例如,按空格字符将每一行文本分割为单词,代码如下:

val words = lines.flatMap(_.split(" "))

在本例中,lines的每一行将被分成多个单词,单词组成的数据流则为一个新的DStream,使用words表示。接下来需要统计单词数量,代码如下:

import org.apache.spark.streaming.StreamingContext._ //计算每一批次中的每一个单词数量
val pairs = words.map(word => (word, 1))
val wordCounts = pairs.reduceByKey(_ + _) 
wordCounts.print() //默认打印前10个元素(DStream中的RDD)到控制台
  1. 启动Spark Streaming

在DStream的创建与转换代码编写完毕后,需要启动Spark Streaming才能真正的开始计算,因此需要在最后添加以下代码:

//开始计算
ssc.start()      
//等待计算结束
ssc.awaitTermination() 

到此,一个简单的单词计数例子就完成了。

05 Spark Streaming数据源

Spark Streaming提供了两种内置的数据源支持:基本数据源和高级数据源。基本数据源是指文件系统、Socket连接等;高级数据源是指Kafka、Flume、Kinesis等数据源。

  1. 文件流

对于从任何与HDFS API兼容的文件系统上的文件中读取数据,可以通过以下方式创建DStream:

streamingContext.fileStream[KeyClass, ValueClass, InputFormatClass](dataDirectory)
  1. Socket流

通过监听Socket端口接收数据,例如以下代码,从本地的9999端口接收数据:

//创建一个本地StreamingContext对象,使用两个执行线程,批处理间隔为1秒
val conf = new SparkConf()
  .setMaster("local[2]")
  .setAppName("NetworkWordCount")

val ssc = new StreamingContext(conf, Seconds(1))
//连接localhost:9999获取数据,转为DStream
val lines = ssc.socketTextStream("localhost", 9999)
  1. RDD队列流

使用streamingContext.queueStream(queueOfRDDs)可以基于RDD队列创建DStream。推入队列的每个RDD将被视为DStream中的一批数据,并像流一样进行处理。这种方式常用于测试Spark Streaming应用程序。

下面以Kafka数据源为例,介绍高级数据源的使用。

首先需要在Maven工程中引入Spark Streaming的API依赖库:

<!--Spark核心库-->
<dependency>
  <groupId>org.apache.spark</groupId>
  <artifactId>spark-core_2.11</artifactId>
  <version>2.4.0</version>
</dependency>

<!--Spark Streaming依赖库-->
<dependency>
  <groupId>org.apache.spark</groupId>
  <artifactId>spark-streaming_2.11</artifactId>
  <version>2.4.0</version>
</dependency>

然后引入Spark Streaming针对Kafka的第三方依赖库(针对Kafka 0.10版本):

<dependency>
  <groupId>org.apache.spark</groupId>
  <artifactId> spark-streaming-kafka-0-10_2.11</artifactId>
  <version>2.4.0</version>
</dependency>

引入所需库后,可以在Spark Streaming应用程序中使用以下格式代码创建输入DStream:

import org.apache.spark.streaming.kafka._
val kafkaStream = KafkaUtils.createStream(streamingContext, 
   [ZK quorum], [consumer group id], [per-topic number of Kafka partitions to consume])

06 DStream操作

与RDD类似,许多普通RDD上可用的操作算子DStream也支持。使用这些算子可以修改输入DStream中的数据,进而创建一个新的DStream。对DStream的操作主要有三种:无状态操作、状态操作、窗口操作。

1、无状态操作

无状态操作指的是,每次都只计算当前时间批次的内容,处理结果不依赖于之前批次的数据,例如每次只计算最近1秒钟时间批次产生的数据。常用的DStream无状态操作算子如表所示。

image-20220501230342715

2、状态操作

状态操作是指,需要把当前时间批次和历史时间批次的数据进行累加计算,即当前时间批次的处理需要使用之前批次的数据或中间结果。使用updateStateByKey()算子可以保留key的状态,并持续不断地用新状态更新之前的状态。使用该算子可以返回一个新的“有状态的”DStream,其中通过对每个key的前一个状态和新状态应用给定的函数来更新每个key的当前状态。

例如,对数据流中的实时单词进行计数,每当接收到新的单词,需要将当前单词数量累加到之前批次的结果中。这里单词的数量就是状态,对单词数量的更新就是状态的更新。定义状态更新函数,实现按批次累加单词数量的代码如下:

/**
 * 定义状态更新函数,按批次累加单词数量
 * @param values 当前批次某个单词的出现次数,相当于Seq(1,1,1)
 * @param state 某个单词上一批次累加的结果,因为可能没有值,所以用Option类型
 */
val updateFunc=(values:Seq[Int],state:Option[Int])=>{

  //累加当前批次某个单词的数量
  val currentCount=values.foldLeft(0)(+) 

  //获取上一批次某个单词的数量,默认值0
  val previousCount= state.getOrElse(0)

  //求和。使用Some表示一定有值,不为None
  Some(currentCount+previousCount)
}

将updateFunc函数作为参数传入updateStateByKey()算子即可对DStream中的单词按批次累加,代码如下:

//更新状态,按批次累加
val result:DStream[(String,Int)]= wordCounts.updateStateByKey(updateFunc)
//默认打印DStream中每个RDD中的前10个元素到控制台
result.print()

3、窗口操作

Spark Streaming提供了窗口计算,允许在滑动窗口(某个时间段内的数据)上进行操作。当窗口在DStream上滑动时,位于窗口内的RDD就会被组合起来,并对其进行操作。

假设批处理时间间隔为1秒,现需要每隔2秒对过去3秒的数据进行计算,此时就需要使用滑动窗口计算,计算过程如图所示(相当于一个窗口在DStream上滑动)。

image-20220501231103471

任何窗口计算都需要指定以下两个参数:

窗口长度:窗口覆盖的流数据的时间长度。必须是批处理时间间隔的倍数。

滑动时间间隔:前一个窗口滑动到后一个窗口所经过的时间长度。必须是批处理时间间隔的倍数。

Spark针对窗口计算提供了相应的算子。例如,每隔10秒计算最后30秒的单词数量,代码如下:

val windowedWordCounts = pairsDStream.reduceByKeyAndWindow(
  (a: Int, b: Int) => (a + b),
  Seconds(30),
  Seconds(10)
)

一些常用的窗口操作算子如表。

image-20220501231232615

4、输出操作

输出操作允许将DStream的数据输出到外部系统,如数据库或文件系统。输出操作触发所有DStream转换操作的实际执行,类似于RDD的行动算子。Spark Streaming定义的输出操作如表。

image-20220501231333840

foreachRDD(func)是一个功能强大的算子,它允许将数据发送到外部系统。理解如何正确有效地使用这个算子非常重要。

5、缓存及持久化

与RDD类似,DStream也允许将流数据持久化到内存中。也就是说,在DStream上使用persist()方法可以将该DStream的每个RDD持久化到内存中。这在DStream中的数据需要被计算多次(例如,对同一数据进行多次操作)时非常有用。对于基于窗口的操作,如reduceByWindow()、reduceByKeyAndWindow(),以及基于状态的操作,如updateStateByKey(),这些都默认开启了persist()。因此,基于窗口操作生成的DStream将自动持久化到内存中,而不需要手动调用persist()。

对于通过网络接收的输入流(如Kafka、Flume、Socket等),默认的持久化存储级别被设置为将数据复制到两个节点,以便容错。

6、检查点

Spark Streaming应用程序必须全天候运行,因此与应用程序逻辑无关的故障(例如,系统故障、JVM崩溃等)不应该对其产生影响。为此,Spark Streaming需要对足够的数据设置检查点,存储到容错系统中,使其能够从故障中恢复。Spark Streaming有两种类型的检查点:元数据检查点、数据检查点。

元数据检查点:元数据主要指配置信息、DStream操作、未完成的批次。

数据检查点:将生成的RDD保存到可靠的存储系统(例如HDFS)中。为了避免系统恢复时间的无限增长,将有状态转换的中间RDD定期存储到可靠系统中,以切断依赖链。

Views: 180

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: 150

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: 135

配置 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: 343

神奇的魔数:0x5f3759df

Quake-III Arena (雷神之锤3)是90年代的经典游戏之一。

该系列的游戏不但画面和内容不错,而且即使计算机配置低,也能极其流畅地运行。这要归功于它3D引擎的开发者约翰-卡马克(John Carmack)。

事实上早在90年代初DOS时代,只要能在PC上搞个小动画都能让人惊叹一番的时候,John Carmack就推出了石破天惊的Castle Wolfstein, 然后再接再励,doom, doomII, Quake...每次都把3-D技术推到极致。他的3D引擎代码资极度高效,几乎是在压榨PC机的每条运算指令。当初MS的Direct3D也得听取他的意见,修改了不少API。

最近,QUAKE的开发商ID SOFTWARE遵守GPL协议,公开了QUAKE-III的原代码,让世人有幸目睹Carmack传奇的3D引擎的原码。

我们知道,越底层的函数,调用越频繁。3D引擎归根到底还是数学运算。

那么找到最底层的数学运算函数(game/code/q_math.c),必然是精心编写的。里面有很多有趣的函数,很多都令人惊奇,估计我们几年时间都学不完。

在/code/game/q_math.c里发现了这样一段代码。

它的作用是将一个数开平方并取倒。

那么你会如何通过代码写出下面这个公式?

f(x) = \frac{1}{\sqrt{x}} \\

如果是我的话,一行代码搞定float y = 1 / sqrt(x)。而《雷神之锤3》实现的方式就非常有趣:

float Q_rsqrt( float number )
{
  long i;
  float x2, y;
  const float threehalfs = 1.5F;

  x2 = number * 0.5F;
  y = number;
  i = * ( long * ) &y; // evil floating point bit level hacking
  i = 0x5f3759df - ( i >> 1 ); // what the fuck?
  y = * ( float * ) &i;
  y = y * ( threehalfs - ( x2 * y * y ) ); // 1st iteration
  // y = y * ( threehalfs - ( x2 * y * y ) ); // 2nd iteration, this can be removed

  #ifndef Q3_VM
  #ifdef __linux__
  assert( !isnan(y) ); // bk010122 - FPE?
  #endif
  #endif
  return y;
}

函数返回1/sqrt(x),这个函数在图像处理中比sqrt(x)更有用。

注意到这个函数只用了一次叠代!(其实就是根本没用叠代,直接运算)。编译,实验,这个函数不仅工作的很好,而且比标准的sqrt()函数快4倍!要知道,编译器自带的函数,可是经过严格仔细的汇编优化的啊!

这个简洁的函数,最核心,也是最让人费解的,就是标注了“what the fuck?”的一句:

i = 0x5f3759df - ( i >> 1 );

再加上

y = y * ( threehalfs - ( x2 * y * y ) );

两句话就完成了开方运算!而且注意到,核心那句是定点移位运算,速度极快!特别在很多没有乘法指令的RISC结构CPU上,这样做是极其高效的。

如果你看到这个0x5f3759df数字,想必有点小懵逼,人生三问开始了:

  1. 它是谁?
  2. 它怎么来的?
  3. 它有什么用?

下面就会介绍这个magic number从何来,去何处。

为何需要这个算法?

当你需要在游戏中实现一些物理效果,比如光影效果、反射效果时,所关注的点其实是某个向量的方向,而不是这个向量的长度,如果能将所有向量给单位化,很多计算就会变得比较简单。

所以在计算中很重要的一点就是计算出单位向量,而在真正在运行代码中需要根据某个向量计算相应的单位向量,根据某个向量(x, y, z)计算法向量公式如下

x_0 = x * \frac{1}{\sqrt{x^2 + y^2 + z^ 2}}  \\
y_0 = y * \frac{1}{\sqrt{x^2 + y^2 + z^ 2}}  \\
z_0 = z * \frac{1}{\sqrt{x^2 + y^2 + z^ 2}}  \\

可以看出来计算一个数的平方根的倒数其实非常频繁,所以需要一个很快的算法去计算平方根倒数。众所周知,乘法和加法在计算机中被设计得非常快,所以x<em>x+y</em>y+z*z计算起来真的非常快,但是

sqrt(xx + yy + zz)求平方根算法很慢,求一个数的倒数即除法也很慢,所以上面的一行代码实现平方根倒数所消耗的时间会特别慢。而Q_rsqrt(x</em>x + y<em>y + z</em>z)里面的代码没有看到任何除法,以及求平方根,里面全部都是乘法、位移等运算速度很快的操作,平方根倒数速算其实是计算的近似数,大约1%的误差,但是运算速度是之前的三倍,下面就会解释这几行代码。

初始化参数

long i;

float x2, y;

const float threehalfs = 1.5F;

x2 = number * 0.5F;

y = number;

首先函数的参数是一个32位浮点数,之后声明一个32位整型变量i,继续声明两个32位浮点数x2,y,声明一个32位浮点数常量threehalfs,即表示\frac{3}{2},之后两行也非常简单,一个是将number的一半赋值给x2,将number赋值给y,你会发现很有趣的一点,就是这些变量不管是整型还是浮点型,其在内存中的长度都是32位,这其实是magic发生的基础。

接着看下面几行的注释(因为直接看代码也看不懂),分别是evil bit hackwhat the fucknewton iteration,这三个注释其实就说明了整个算法重要的三步,在真正解释这三个步骤之前,先来说说浮点数的在内存中表示。

二进制数的表示方法

在日常生活中,我们通常以十进制的方式表示现实生活中的各种数,这种数被称为真值,对于负数,我们会在数的前面加一个‘-’号表示这个数是小于0,‘+’号(通常不写)去表示一个正数。而在计算机中只能表示0、1两种状态,所以下正负号在计算机中会以0、1的形式表示,通常放在最高位,作为符号位。为了能够方便对这些机器数进行算术运算、提高运算速度,计算机设计了许多种表示方式,其中比较常见的是:原码、反码、补码以及移码。下面主要以8bit长度的数0000 0100,即十进制数4,去介绍这几种表示方式

  • 原码

原码表示方式是最容易理解的,首先第一位为符号位,后面七位表示的就是真值,如果表示负数,只需要将符号位置为1即可,后面七位依然为真值,所以4与-4的原码为0000 0100 和 1000 0100。所以很容易就得出原码的表示范围为[-127, 127],会存在两个特殊的数为+0-0

  • 反码

正数的反码即是其原码,而负数的反码就是在保留符号位的基础上,其他位全部取反,所以4与-4的反码为0000 0100 和 1111 1011。所以反码表示范围为[-127, 127],依然存在两个特殊的数为+0-0

  • 补码

正数的补码即是其原码,而负数的补码就是在保留符号位的基础上,其他位全部取反,最后加1,即在反码的基础上+1,所以4与-4的补码是0000 0100 和 1111 1100。所以补码表示范围为[-128,127],之前的-0在补码中被表示成了-128,可以多表示一个数,另外补码统一了正数0和负数0。

  • 移码

假设数值用八位有符号二进制表示, 则移码表示在数轴上将真值在范围从[-128, 127]变为[0, 255],即向上偏移量128,即将补码的符号位取反(将负数平移到正数区间内),所以4与-4的移码为1000 0100 和 0111 1100,移码的好处是可以直观方便的比较两数在数轴上的大小。

浮点数

先考虑一个问题,如果你用32位二进制如何表示4.25?可能会是这样:

0000 0000 0000 0100 . 0100 0000 0000 0000

这放在普通十进制,这种想法其实非常常见,但是这种方式放在二进制世界中,总共1位符号位,15号整数位,16位小数位,总共表示数的范围只有[-2^{15}, 2^{15}],而对于长整型的范围却有[-2^{31}, 2^{31}],差的倍数有6w5,可见这种方式为了小数表示抛弃了一半的位数,得不偿失,所以有人提出了IEEE754标准。

IEEE754

在描述这个标准前,先在这里说下科学计数法,在十进制中科学计数法表示如下

\begin{aligned} & 12300 && 1.23 * 10^{4} \\ & 0.00123 && 1.23 * 10^{-3} \\ \end{aligned}  \\

同样的,可以将科学计数法运用到二进制中

\begin{aligned} & 11,000 && 1.1 * 2^{4} \\ & 0.0101 && 1.01 * 2^{-2} \\ \end{aligned} \\

所以IEEE754也是采用的是科学计数法的形式,会将32位数分为以下三部分

  • Sign Bit

首先第一位是符号位,0表示正数,1表示负数,而在平方根的计算中,明显不会涉及到负数,所以第一位肯定是为0的。

  • Exponent

第二部分用8位bit表示指数部分(乘上2的几次方),可以表示数的范围是[0, 255],但是这个只能表示正数,所以需要把负数也加进来,而IEEE754标准中阶码表示方式为移码,之所以要表示为移码的方式是在浮点数比较中,比较阶码的大小会变得非常简单,按位比较即可。不过和正常的移码有一点小区别是,0000 00001111 1111用来表示非规格化数与一些特殊数,所以偏移量从128变为127,表示范围也就变成了[-127, 126]

举个例子,众所周知啊,4这个数的8bit真值为0000 0100,加上127的偏移量变成131,即阶码的二进制表示为1000 0011

  • Mantissa

二进制科学计数法中小数部分就是剩余的23位,由于第一位二进制数总是1,所以用这23位bit表示剩余的数值,真正计算的时候再在最高位加上1即可。

下面我们以9.625这个数改成IEEE 754标准的机器数:

首先变为相应的二进制数为1001.101,用规范的浮点数表达应为1.001101 * 2^{3},所以符号段为0,指数部分的移码为1000 0010,有效数字去掉1后为001101,所以最终结果为

0 |10000010 |00110100000000000000000

二进制数转换

在前面我们介绍了一个浮点数是如何在计算机中以机器数来表示的,现在我们要对这个浮点数进行一些骚操作,以方便之后对这个浮点数的处理。

同样的,我们以9.625这个数为例,首先我们令阶码的真值为E,则有

E = 1000,0010_{(2)} = 130_{(10)}  \\

令余数为M,则有

M = 00110100000000000000000_{(2)} = {0.001101 * 2^{23}}_{(2)} \\

现在我们先不认为这个是浮点数,这个数就是32位长整型,令这个数为L,它所表示的十进制数便是这样

L = 2^{23} * E + M \\

这个数L便是这32位所表示的无符号整型,这个数后面有大用,我们先暂且把它放在这儿,后续再来看。

然后我们通过一个公式来表示这个浮点数F的十进制数

F = (1 + \frac{M}{2^{23}}) * 2^{E - 127} \\

尾数加1是因为在IEEE754标准中把首位的一去掉了,所以计算的时候需要把这个一给加上,然后阶码减去127是因为偏移为127,要在8位真值的基础上减去127才是其表示真正的值。

然后有趣的事情就发生了,我们现在将这个数F取下对数

\log_2 F = \log_2 (1 + \frac{M}{2^{23}}) + E - 127 \\

在这个公式,E-127自然非常好计算,但是前面的对数计算起来是比较麻烦的,所以我们可以找个近似的函数去代替对数。

现在看下y=log_2(1 + x)y=x的图像

根据图像,我们很容易得出下面这个结论

\lim_{x \to 0^+} \log_2 (1 + x) = x  \\
lim_{x \to 1^-} \log_2 (1 + x) = x \\

然后很容易发现,在[0, 1]这个范围内,\log_2(x+1)x其实是非常相近的,那样我们可以取一个在[0,1]的校正系数\mu$,使得下面公式成立

\log_2 (x + 1) \approx x + \mu \\

到此,我们知道了怎样去简化对数,所以我们可以将这个简化代入上面浮点数表示中,就可以得到

\begin{aligned} \log_2 F &= \log_2 (1 + \frac{M}{2^{23}}) + E - 127 \\ & \approx (\frac{M}{2^{23}} + \mu) + E - 127 \\ & = (\frac{M}{2^{23}} + \frac{E * 2^{23}}{2^{23}}) + \mu - 127 \\ & = \frac{1}{2^{23}} (M + E * 2^{23}) + \mu - 127 \end{aligned} \\

看见这个M + E * 2 ^{23}\\应该比较熟悉吧,这个数就是浮点数F的二进制数L,然后代入约等式中得到

\log_2 F \approx \frac{L}{2^{23}} + \mu - 127 \\

在某种程度上,不考虑放缩与变换,我们可以认为浮点数的二进制表示L其实就是其本身F的对数形式,即

\log_2 F \approx L \\

三部曲

经过上面一系列复杂的数字处理操作,我们终于可以开始我们的算法三部曲了

evil bit hack

众所周知啊,每个变量都有自己的地址,程序运行的时候就会通过这个地址拿到这个变量的值,然后进行一系列的计算,比如i和y在内存中会这样表示

我们对长整型这些数据很好进行位移运算,比如我想将这个数乘以或者除以2^n\\,只需要左移或者右移N个位就可以,但是浮点数明显无法进行位运算,它本身二进制表示就不是为了位运算设计的。

然后,现在就会提出一个想法,我把float转成int,然后进行位运算不就行了,代码如下long i = (long) y

假设y为3.33,进行长整型强转后,C语言会直接丢弃尾数,i也就变成了3,丢失这么多精度,谁干啊,如果我们想一位都不动地进行位运算,就是下面这份代码i = <em> ( long </em> ) &y

这个强转过程其实没改变内存中任何东西,首先它并没有改变0x3d6f3d791这个地址,也没改变这个地址中所存储的数据,可以认为改变的是C语言的“认知”,原本是要以IEEE754标准去读取这个地址中数据,但是C语言现在认为这个是长整型的地址,按长整型方式读取就行。所以(long<em>)&y代表0x3d6f3d791这个地址中存储的是长整型数据,然后我通过</em>运算符从这个地址中拿出数据,赋值给i这个长整型变量。

这行代码做了什么事呢,&y首先将浮点数y的地址找出来,可以认为就是0x3d6f3d79这个地址,它的类型其实是float <em>,C语言便会以浮点数的形式将这个数取出来,而想让C语言认为这个是长整型类型,就必须进行地址的类型的强转,将float </em>强转成long *

这样我们其实避开了数字本身的意义,而是通过地址变换完整地拿出了这个二进制数据。

what the fuck

众所周知啊,位运算中的左移和右移一位分别会使原数乘以或者除以2,比如

\begin{aligned} x & & 110 &= 6 \\ x << 1 & & 1100 &= 12 \\ x >> 1 & & 11 &= 3 \\ \end{aligned} \\

我们想办法把平方根倒数做一个简单的转化

\frac{1}{\sqrt x} = x^{-\frac{1}{2}} \\

这个等式其实就是我们的最终目标,接下的计算就会逐渐往这个等式靠近,得到一个近似的结果。

在前面一步的叙述中,我们得出这样一个结论,浮点数的二进制表示其实就是其本身的对数形式,想要求的浮点数存储在y中,则有

y = 9.625 \\  \\
log_2 (y) \approx i = 01000001 \ 00011010 \ 00000000 \ 00000000 \\

也就是说i中其实存储着y的对数,当然还需要进行一系列的转换与缩放。

前面提到过,直接运算一个数的平方根倒数,所以不如直接计算平方根倒数的对数,然后就会有如下等式

\begin{aligned} \log_2 (\frac{1}{\sqrt y}) = \log_2 (y^{- \frac{1}{2}}) & = - \frac{1}{2} \log_2(y) \\ \log_2 (y) & \approx i \\ \log_2 (\frac{1}{\sqrt y}) & \approx -(i >> 1) \end{aligned}  \\

除法同样计算速度是比较慢的,所以我们用右移代替除法,这也就解释了-(i >> 1)其实是为了计算\log_2 (\frac{1}{\sqrt y})的结果。

所以0x5f3759df这个数到底咋来的,为啥要这么计算,并且-(i >> 1)并不完全是\log_2 (\frac{1}{\sqrt y})的近似值啊,根据公式还得除以2^{23},还得加上一定的误差。

先别急,我们先算下这个magic number是咋来的,令\Gamma为y的平方根倒数,则有

\begin{aligned} \Gamma & = \frac{1}{\sqrt y} \\ \log_2 \Gamma & = \log_2 (\frac{1}{\sqrt y}) = -\frac{1}{2} \log_2 (y) \end{aligned} \\

然后我们代入上面那个浮点数对数的公式中,则有

\begin{aligned} \log_2 \Gamma & = \frac{1}{2^{23}} (M_{\Gamma} + E_{\Gamma} * 2^{23}) + \mu - 127 \\ \log_2 y & = \frac{1}{2^{23}} (M_y + E_y * 2^{23}) + \mu - 127 \\ \frac{1}{2^{23}} (M_{\Gamma} + E_{\Gamma} * 2^{23}) + \mu - 127 & = -\frac{1}{2}(\frac{1}{2^{23}} (M_y + E_y * 2^{23}) + \mu - 127) \end{aligned} \\

现在经过一定的变换能够得到下面这个式子

\begin{aligned} (M_{\Gamma} + E_{\Gamma} * 2^{23}) & = \frac{3}{2}2^{23}(127 - \mu) -\frac{1}{2} (M_y + E_y * 2^{23}) \\ & = 0x5f3759df - (i >> 1) \end{aligned} \\

这里\mu是就是之前简化对数计算引进的误差,通过一定计算,得到最合适的\mu,来得到\Gamma的近似值。这个计算过程偏向纯数学化的,具体过程请查阅《FAST INVERSE SQUARE ROOT》这篇论文。

原算法中取的\mu值为0.0450465,然后计算一下

\frac{3}{2}2^{23}(127 - \mu) = \frac{3}{2}2^{23}(127 - 0.0450465) = 1597463007 \\

这个数的十六进制就是0x5f3759df,也就是上面提到的那个magin number,然后我们根据这个数去计算得到的近似值,并与真正的\frac{1}{\sqrt y}进行比较,具体函数图像如下

从上面的图像可以看到,在[1,100]这个区间内,所得到的近似值曲线已经和原始值制拟合的比较好了,这样我们已经完成了前面几个比较重要的步骤。y = <em> ( float </em> ) &i;

这样我们通过和evial bit hack的逆向步骤,即将一个长整型的内存地址,转变成一个浮点型的内存地址,然后根据IEEE754标准取出这个浮点数,即我们要求的\Gamma的近似值,其实到这里应该是算法差不多结束了,但是这个近似值还存在一定的误差,还需要经过一定的处理降低误差,更接近真实值。

Newton Iteration

本身我们已经得到了一个比较好的近似值,但是仍然存在一定的误差,而牛顿迭代法可以这个近似值更加接近真实的值,近一步减少误差。

牛顿迭代法本身是为了找到一个方程根的方法,比如现在有一个方程f(x)=0,需要找到这个方程的根,但是解方程嘛,是不会解方程的,所以可以找到一个近似值来代替这个真正的解。

如上图,假设f(x)=0的解为x=x_c,我们需要首先给一个x的近似值x_0,通过这个x_0,不断求得一个与x_c更接近的值。

(x_0, f(x_0))处做切线,切线的斜率就是f(x)x=x_0的导数f'(x_0),然后来求这条切线与x轴的交点x_1,则有

x_1 = x_0 - \frac{f(x_0)}{f'(x_0)}  \\

这样我们就完成了一次迭代,从图像上可以看见x_1x_0更接近于真正的解,下一次迭代基于x_1进行同样的步骤,就能得到比x_1更好的近似值x_2,所以牛顿迭代公式非常简单

x_{n + 1} = x_n - \frac{f(x_n)}{f'(x_0)} \\

当这个迭代次数接近无限时,x_n也就越接近真正的解。

而最后一行就是经过一次迭代后的简化公式,这个公式怎么来的呢。对于一个浮点数a,要求它的平分根倒数,则有

y = \frac{1}{\sqrt a} \\

通过这个公式能构成一个函数

f(y) = \frac{1}{y^2} - a \\

求这个y值,其实就是求f(y)=0的根,所以迭代公式就是

\begin{aligned} y_{n+1} & = y_n - \frac{f(y_n)}{f'(y_n)} \\ & = y_n - \frac{y_n^{-2} - a}{-2 * y_n^{-3}} \\ & = y_n * (\frac{3}{2} - \frac{1}{2}* a * y_n^2) \end{aligned} \\

这个公式对应的算法中的代码y = y <em> (threehalfs - ( x2 </em> y * y ) )

至此,你的代码就更接近真实的y=\frac{1}{\sqrt x},在更接近真实答案的同时,运行速率也大大提升。仅仅牺牲了一点点的准确性,却能提高整个的速度,这其实就是算法在优化中的一个比较重要的点。

Lomont -“Fast Inverse Square Root”。

普渡大学的数学家Chris Lomont看了这个算法以后觉得有趣,决定要研究一下卡马克弄出来的这个猜测值有什么奥秘。Lomont也是个牛人,在精心研究之后从理论上也推导出一个最佳猜测值,和卡马克的数字非常接近, 0x5f37642f。卡马克真牛,他是外星人吗?

传奇并没有在这里结束。Lomont计算出结果以后非常满意,于是拿自己计算出的起始值和卡马克的神秘数字做比赛,看看谁的数字能够更快更精确的求得平方根。结果是卡马克赢了... 谁也不知道卡马克是怎么找到这个数字的。

最后Lomont怒了,采用暴力方法一个数字一个数字试过来,终于找到一个比卡马克数字要好上那么一丁点的数字,虽然实际上这两个数字所产生的结果非常近似,这个暴力得出的数字是0x5f375a86。

Lomont为此写下一篇论文,"Fast Inverse Square Root"。

给出了InvSqrt函数的定义:

float InvSqrt(float x)
{
  float xhalf = 0.5f * x;
  int i = *(int *)&x;
  i = 0x5f3759df - (i>>1);
  x = *(float *)&i;
  x = x * (1.5f - xhalf * x * x);
  return x;
}

本文参考自:

知乎Elb(字节跳动大力教育前端团队)-https://wjrsbu.smartapps.cn/zhihu/article?id=445813662&isShared=1&_swebfr=1&_swebFromHost=baiduboxapp

Views: 192