Spark Streaming(5):Spark Window function in Java

首先,看下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}
点赞
收藏

评论区

加载中...

相关推荐

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

一篇文章带你了解JavaScript时间

一、前言setTimeout(function,milliseconds)在等待指定的毫秒数后执行函数。setInterval(function,milliseconds)setTimeout()相同,但会重复执行。二、时间事件窗口对象允许在指定的时间间隔执行代码。时间间隔称为定时事件。1\.setTimeout()方法window.set

Java日期时间API系列31

  时间戳是指格林威治时间1970年01月01日00时00分00秒起至现在的总毫秒数,是所有时间的基础,其他时间可以通过时间戳转换得到。Java中本来已经有相关获取时间戳的方法,Java8后增加新的类Instant等专用于处理时间戳问题。 1获取时间戳的方法和性能对比1.1获取时间戳方法Java8以前

java——20171121

!(http://a.51jsoft.com/uploads/default/original/1X/c542896b094a42a5653fb75adf6cdacd6e35d12e.png)!(https://static.oschina.net/uploads/space/2017/1121/210719_G80Z_3715033.png)

个推分享Spark性能调优指南:性能提升60%↑ 成本降低50%↓

前言Spark是目前主流的大数据计算引擎,功能涵盖了大数据领域的离线批处理、SQL类处理、流式/实时计算、机器学习、图计算等各种不同类型的计算操作,应用范围与前景非常广泛。作为一种内存计算框架,Spark运算速度快,并能够满足UDF、大小表Join、多路输出等多样化的数据计算和处理需求。作为国内专业的数据智能服务商,个推从早期的1.3版本便引入Spark,

SparkSql学习1 —— 借助SQlite数据库分析2000万数据

总所周知,Spark在内存计算领域非常强势,是未来计算的方向。Spark支持类Sql的语法,方便我们对DataFrame的数据进行统计操作。但是,作为初学者,我们今天暂且不讨论Spark的用法。我给自己提出了一个有意思的思维游戏:Java里面的随机数算法真的是随机的吗?好,思路如下:1\.取样,利用Java代码随机生成2000万条01