
1、Spark Streaming的三种操作
也可称为三种应用场景
1:无状态操作。如操作某一个RDD。
操作某一段时间的数据。如设置为5秒,则操作这5秒之内的数据就是指无状态的操作。
2:状态操作。如操作连续的一组RDD。
是不是指从开始到目前所有RDD的数据,如updateStateByKey就可以称为有状态的操作。
在进行updateStateByKey时,必须要之前设置了checkpoint目录。因为每一次处理都包含前面的所有的数据RDD,此时就必须要将之前的数据先进行缓存。而checkpoint类似于缓存点。这样就减少了内存的使用。
3:window操作,即窗口操作。
即一段时间的数据。这个时间可以跨越多个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
共同交流,共同进步。欢迎指正。




