暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

SparkStreaming

Coding On Road 2018-01-13
475


1、Spark Streaming的三种操作


也可称为三种应用场景

1:无状态操作。如操作某一个RDD

操作某一段时间的数据。如设置为5秒,则操作这5秒之内的数据就是指无状态的操作。

2:状态操作。如操作连续的一组RDD

是不是指从开始到目前所有RDD的数据,如updateStateByKey就可以称为有状态的操作。

在进行updateStateByKey时,必须要之前设置了checkpoint目录。因为每一次处理都包含前面的所有的数据RDD,此时就必须要将之前的数据先进行缓存。而checkpoint类似于缓存点。这样就减少了内存的使用。

3window操作,即窗口操作。

即一段时间的数据。这个时间可以跨越多个5(假设设置为5秒的话)

2、SparkStreaming的编程模型

Spark Streaming的编程模型为 DStream - 离散数据流 - 由很多RDD组成。

Discretized Stream表示为数据流,实现为 RDD 序列,它的创建方式:

1. 从流输入源创建

2. 从现有 DStreams

3. 通过 transformation 操作转换而来


以下是从外部文件读取文件的操作。


3、textFileStream读取目录下的新增文件

注意:

windows上使用textFileStream新文件无论如何测试也没有读取到里面的内容,只能读取到文件,应该与windows管理文件的方式有关。

所以,请在linux环境下进行测试。

开以下代码:

package cn.spark.streaming
import org.apache.log4j.{Level, Logger}
import org.apache.spark.SparkConf
import org.apache.spark.streaming.dstream.DStream
import org.apache.spark.streaming.{Seconds, StreamingContext}
/**
  * 读取某个目录下的文件<br>
  * 如果有新的文件增加,将会读取这个文件形成DStream
  */
object Demo02_FileDirectory {
  def main(args: Array[String]): Unit = {
    if (args.length < 2) {
      println("Usage:参数1读取目录 参数2输出目录");
      return;
    }
    //设置日志的级别
    Logger.getLogger("org").setLevel(Level.WARN);
    //必须要大于两个线程,因为一个进行监听,一个用于处理数据
    val conf: SparkConf = new SparkConf().setAppName("StreamingTextFile");
    val ssc: StreamingContext = new StreamingContext(conf, Seconds(5));
    //设置每5秒读取某个目录下的文件,如果这个目录不存在,则将会抛出异常,所以,请提前创建这个目录
    //注意:监听windows的目录将不成功
    //请在Linux下测试
    val lines: DStream[String] = ssc.textFileStream(args(0));
    val ds: DStream[String] = lines.flatMap(line => line.split("\\s+"))
      .map((_, 1))
      .reduceByKey(_ + _).map(kv => kv._1 + "\t" + kv._2);
    //输出前10
    ds.print();
    ds.saveAsTextFiles(args(1));
    //启动
    ssc.start();
    //等待用户停止
    ssc.awaitTermination();
  }
}

 

将上面的代码进行打包如:streaming.jar

 

使用spark-submit提交:
$ spark-submit --class cn.spark.streaming.Demo02_FileDirectory --master local[2] streaming.jar home/wangjian/a1/ home/wangjian/a2

第一个目录为读取的文件目录,第第二个目录为输出的目录:

则输出数据以后,目录的格式为:

[wangjian@hadoop201 a2]$ tree

.

├── a2-1515734710000

│   ├── part-00000

│   ├── part-00001

│   └── _SUCCESS

├── a2-1515734715000

│   ├── part-00000

│   ├── part-00001

│   └── _SUCCESS

 

现在在 a1目录下,通过vim创建一个,且写入一些数据。

再通过cp/mv将一个文件放到a1目录下,检查控制台的输出:

-------------------------------------------

Time: 1515734750000 ms

-------------------------------------------

Hello   1

Mary    2

Rose    1

Jack    2

 

然后再读取a2目录下的所有文件:

[wangjian@hadoop201 a2]$ cat a2-*/*

Hello   1

Mary    2

Rose    1

Jack    2

 

可见,可以读取到所有的数据。

 

1、也可以将数据保存到hdfs

只要将输出目录设置成hdfs://..即可。

$ spark-submit --class cn.spark.streaming.Demo02_FileDirectory --master local[2] streaming.jar home/wangjian/a1 hdfs://hadoop201:8020/out/out01

 

4、fileStream处理文件过虑

1读取hdfs目录变化

注意在监控hdfs上目录文件的变化时,必须要对_COPYING_的文件进行过虑掉,因为HDFS在上传文件时,会在当前目录下形成一个xxx.txt_COPYING_的临时文件,待上传完成以后,此文件就会删除,如果对这个文件,不做处理在删除时,SparkStreaming将会抛出一个异常:FileNotfoundException。所以,不能直接使用textFileStream直接监控hdfs文件系统的某个目录,必须要使用fileStream。它的第二个参数filter=>(path:org.apache.hadoop.fs.Path)可以指定自己的过虑方式。

以下是完整代码:

 

package cn.spark.streaming
import org.apache.hadoop.fs.Path
import org.apache.hadoop.io.{LongWritable, Text}
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat
import org.apache.log4j.{Level, Logger}
import org.apache.spark.SparkConf
import org.apache.spark.streaming.dstream.{DStream, FileInputDStream, InputDStream}
import org.apache.spark.streaming.{Seconds, StreamingContext}
/**
  * 由于在将文件上传到hdfs时,hdfs上会创建一个临时文件:_COPYING<br>
  * 现在要对这个文件忽略不读取
  */
object Demo03_fileStream {
  def main(args: Array[String]): Unit = {
    //设置日志的级别
    Logger.getLogger("org").setLevel(Level.WARN);
    //必须要大于两个线程,因为一个进行监听,一个用于处理数据
    val conf: SparkConf = new SparkConf().setAppName("StreamingTextFile");
    val ssc: StreamingContext = new StreamingContext(conf, Seconds(5));
    //设置文件过虑
    val lines: InputDStream[(LongWritable, Text)] =
      ssc.fileStream[LongWritable, Text, TextInputFormat](directory = "hdfs://hadoop201:8020/a1", //读取的目录
        filter = ((path: Path) => (!path.getName.contains("_COPYING_"))), //文件过虑,这儿必须要指定参数的类型
        newFilesOnly = true);
    //注意处理方式与直接读取textFileStream的小小区别
    val ds: DStream[String] = lines.map(kv => kv._2)//先获取行数据,不获取行偏移量
      .flatMap(txt => txt.toString.split("\\s+"))
      .map((_, 1)).reduceByKey(_ + _)
      .map(kv => kv._1 + "\t" + kv._2);
    //只输出前10
    ds.print();
    //声明保存的目录
    ds.saveAsTextFiles("hdfs://hadoop201:8020/out/out002");
    ssc.start();
    ssc.awaitTermination();
  }
}

 

测试结果,正常:

-------------------------------------------

Time: 1515739195000 ms

-------------------------------------------

Hello   1

Alex    1

Mike    1

Rose    1

Jack    1

 

查看hdfs目录上的数据:

[wangjian@hadoop201 ~]$ hdfs dfs -ls /out/

Found 33 items

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:55 /out/out002-1515740135000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:55 /out/out002-1515740140000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:55 /out/out002-1515740145000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:55 /out/out002-1515740150000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:55 /out/out002-1515740155000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:56 /out/out002-1515740160000

drwxr-xr-x   - wangjian supergroup          0 2018-01-12 14:56 /out/out002-1515740165000

 

共同交流,共同进步。欢迎指正。

文章转载自Coding On Road,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论