首先,看下window函数的图解:

下面这个代码是计算一分钟之内的单词数量统计,每两秒获取一次数据,同时处理数据时间也是两秒,窗口大小为1分钟
1.数据源
1package com.ssm.test; 2 3import java.io.BufferedReader; 4import java.io.InputStreamReader; 5import java.io.PrintWriter; 6import java.net.ServerSocket; 7import java.net.Socket; 8 9public class SocketTest { 10 11 public static void main(String[] args) { 12 try{ 13 ServerSocket server=new ServerSocket(9999); 14 System.out.println("The Server has started"); 15 16 Socket socket=server.accept(); 17 System.out.println("There was a socket client connected"); 18 19 BufferedReader br=new BufferedReader(new InputStreamReader(System.in)); 20 String line=br.readLine(); 21 PrintWriter writer=new PrintWriter(socket.getOutputStream()); 22 23 while(!line.equals("end")){ 24 writer = new PrintWriter(socket.getOutputStream()); 25 writer.println("If you are not brave enough, no one will back you up."); 26 writer.flush(); 27 System.out.println(System.currentTimeMillis() + "\nServer send:\t"+line); 28 Thread.sleep(2000); 29 } 30 writer.close(); 31 socket.close(); 32 server.close(); 33 }catch(Exception e) { 34 e.printStackTrace(); 35 } 36 } 37}
2.spark streaming处理代码
1package com.ssm.test; 2 3import java.util.Arrays; 4import java.util.Iterator; 5 6import org.apache.spark.SparkConf; 7import org.apache.spark.api.java.function.FlatMapFunction; 8import org.apache.spark.api.java.function.Function2; 9import org.apache.spark.api.java.function.PairFunction; 10import org.apache.spark.streaming.Duration; 11import org.apache.spark.streaming.Durations; 12import org.apache.spark.streaming.api.java.JavaDStream; 13import org.apache.spark.streaming.api.java.JavaStreamingContext; 14 15import scala.Tuple2; 16 17public class WordsCountTest { 18 19 public static void main(String[] args) throws Exception { 20 21 SparkConf conf = new SparkConf().setMaster("local[2]").setAppName("WordsCount"); 22 JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(2)); 23 JavaDStream<String> lines = jssc.socketTextStream("127.0.0.1", 9999).window(new Duration(60000)); 24 25 lines.flatMap(new FlatMapFunction<String, String>(){ 26 public Iterator<String> call(String x){ 27 return Arrays.asList(x.split(" ")).iterator(); 28 } 29 }).mapToPair(new PairFunction<String, String, Integer>(){ 30 public Tuple2<String, Integer> call(String s){ 31 return new Tuple2<String, Integer>(s, 1); 32 } 33 }).reduceByKey(new Function2<Integer, Integer, Integer>(){ 34 public Integer call(Integer i1, Integer i2){ 35 return i1+i2; 36 } 37 }).print(); 38 39 jssc.start(); 40 jssc.awaitTermination(); 41 42 } 43 44}