ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Java面试——Spark原理及应用(二)

Java面试——Spark原理及应用(二) Spark原理及应用2.3、Spark Streaming的使用2.3.1、Spark Streaming的介绍2.3.2、创建一个Spark Streaming应用2.3.3、DStream和RDD的关系2.3.4、Spark Streaming的数据源2.3.5、DStream的操作2.3.6、DStream的数据持久化2.3.7、Spark Streaming的性能优化2.4、Spark SQL、DataFrame、DataSet的使用2.4.1、Spark SQL简介2.4.2、Spark SQL查询语句2.4.3、DataFrame2.4.4、DataSet2.4.5、创建一个Spark SQL应用2.4.6、Spark SQL的视图操作2.4.7、创建DataSet2.4.8、DataSet与RDD相互转换2.4.9、DataFrame数据的加载2.4.10、DataFrame数据的保存2.5、Spark Structured Streaming的使用2.5.1、Spark Structured Streaming简介2.5.2、Spark Structured Streaming的数据模型2.5.3、创建一个Spark Structured Streaming应用2.3、Spark Streaming的使用2.3.1、Spark Streaming的介绍Spark Streaming基于Spark API的流式计算扩展它实现了一个高吞吐量、高容错的流式计算引擎。SparkStreaming从多种数据源获取数据如Kafka、Flume、Kinesis、TCP等​在获取数据后使用高级函数如map、reduce、join、window组成的计算逻辑单元进行数据处理最终将处理后的数据实时推送到消息服务、文件系统、数据库等。除了基本的流式计算应用程序还可以在数据流上应用Spark的机器学习和图形处理算法。Spark Streaming的流程如图所示。Spark Streaming接收实时数据流并将数据分成很小的Batch然后将这些Batch交给Spark Engine以微批量的形式进行实时处理生成最终结果流。Spark Streaming的处理过程如图所示。在API层面Spark Streaming将实时数据流封装为DStreamDiscretized Stream的高级抽象以表示连续的数据流。DStream可以从Kafka、Flume和Kinesis等实时数据流创建也可以通过在其他DStream上应用高级操作来创建。在内部DStream表示为一系列RDD。SparkStreaming内部流计算的过程即Spark RDD的转化过程。2.3.2、创建一个Spark Streaming应用在之前的Spark项目基础上创建Spark Streaming项目具体步骤如下。1添加Maven依赖。打开pom.xml文件并将Spark Streaming依赖加入项目中具体代码如下。** 2新建Streaming类。**在项目的Java目录下新建一个名为Streaming的Java类并在类中输入以下代码构建一个简单的Spark Streaming程序。上述代码定义了一个名为NetworkWordCount的SparkStreaming计算引擎该程序实时监听Localhost的9999端口。当有数据输入时Spark Streaming实时接收9999端口上的数据并计算词频最终将结果输出。在SparkStreaming应用的创建过程中包括以下核心步骤。①定义JavaStreamingContext定义JavaStreamingContext需要传入SparkConf集群配置信息在创建的时候可以定义Stream的持续时间即SparkStreaming多久运算一次实时流数据这里设置为1s一般意义上可理解为数据延迟1s计算。具体代码如下。②定义实时数据源Spark Streaming中的数据源可以是TCP端口上的实时数据、Kafka中的数据、文件目录中的数据等通过调用JavaStreamingContext的转换函数将 其转换为Spark的实时数据流为后续流计算提供数据源。数据源通过JavaStreamingContext定义如下代码通过创建一个名为lines的JavaReceiverInputDStream表示一个监听了Localhost的9999端口的TCP数据源。③定义DStream转换操作DStream表示从数据源接收的实时数据流。如下代码中实时流为9999端口上用户输入的每一行文本。首先通过调用flatMap操作按照空格将每行文档分割为单词该过程被称为数据源到JavaDStream的转换。然后调用words.mapToPair​将DStream进一步转化为键值对的DStream。接着通过pairs.reduceByKey​获得每批数据中的单词频率。最后通过wordCounts.print​方法打印每秒计算的词频计数。具体代码如下。④定义DStream结果输出上述代码通过wordCounts.print​方法将结果打印当然也可以将计算结果存储在数据库中发送到消息队列或者存储到HDFS等文件系统中。⑤启动Spark StreamingSpark Streaming的启动通过调用JavaStreamingContext的start​方法实现。3启动TCP服务端。在Linux上通过Netcat启动一个TCP服务端具体代码如下。4打包和作业提交。在项目的根目录下输入如下命令对项目进行打包打包后的程序在Target目录下。Spark Streaming的作业提交与一般的作业提交类似具体代码如下。在作业提交后通过http://ip:4040/jobs/可以查看启动的作业如图所示。在Netcat窗口输入“hello java hello spark”​查看Spark日志看到有以下统计结果输出。2.3.3、DStream和RDD的关系DStream是Spark Streaming提供的基本抽象它表示连续的数据流可以是从其他数据源接收的输入数据流也可以是通过转换输入流生成的。在Spark内部DStream由一系列连续的RDD表示这是Spark对不可变分布式数据集的抽象。应用于DStream的任何操作都转换为底层RDD上的操作。这些底层RDD转换由Spark计算引擎完成。DStream操作隐藏了大部分细节并为开发人员提供了更高级别的API以方便使用。DStream的计算流程和RDD的转换关系如图所示。2.3.4、Spark Streaming的数据源Spark Streaming提供如下两种内置流式数据的数据源。基本来源StreamingContext API中直接提供的源如文件系统和Socket连接等。高级资源Kafka、Flume、Kinesis等数据源可通过扩展实现其他依赖。如果要在应用程序中并行接收多个数据流则可以通过定义多个Receiver来实现这些接收器将同时接收多个数据源上的数据。需要注意的是分配给Spark Streaming应用程序的核心数必须大于接收器的数量因为系统在接收数据的同时需要有足够的线程资源来处理数据。1常用的基本数据源。①SocketSpark启动一个常驻内存的线程来实时监听Socket上的数据变化。具体通过如下方法创建DStream。②HDFS文件流从与HDFS API兼容的任何文件系统HDFS、S3、NFS上读取数据。具体通过如下方法创建DStream。③简单文件流从简单的文件目录上读取数据。2常用的高级数据源。①Flume利用Flume Spark Streaming从Flume中获取数据。②Kafka利用Kafka Spark Streaming从Kafka中获取数据。③Kinesis利用Kinesis Spark Streaming从Kinesis中获取数据。2.3.5、DStream的操作Dstream的操作可以分为普通转换操作、窗口转换操作和输出操作等。1常用的普通转换操作。①mapfunc​原DStream的每个元素都通过func函数返回一个新的DStream。②flatMapfunc​与map操作类似不同的是每个输入元素都可以被映射出0个或多个输出元素。③filterfunc​在原DStream上过滤出func函数返回仅为true的DStream。④repartitionnumPartitions​设置DStream的分区大小。⑤unionotherStream​将2个DStream合并。⑥count​​对原DStream内部所含有的RDD的元素数量进行统计。⑦reducefunc​使用函数func将原DStream中每个RDD的元素都进行聚合操作。⑧countByValue​​计算DStream中每个RDD内的元素出现的频次。⑨reduceByKeyfunc​[numTasks]​​根据DStream的Key进行聚合操作。⑩joinotherStream​[numTasks]​​当被调用的类型分别为KeyValue1和KeyValue2键值对的2个DStream时返回类型为KeyValue1Value2键值对的一个新DStream。⑪cogroupotherStream​[numTasks]​​当被调用的两个DStream分别含有KeyValue1和KeyValue2键值对时返回一个KeySeq[Value1]​Seq[Value2]​类型的新DStream。2窗口转换操作。要理解Spark的窗口转换操作首先要理解批处理间隔、窗口间隔、滑动间隔等基础概念。①批处理间隔Spark Streaming中的数据处理是按批进行的而数据采集是实时逐条进行的。Spark Streaming内部设置了批处理间隔Batch Duration​Spark会定时把该批处理间隔内接收的数据汇总起来新组成一个微批数据Batch Data并提交到Spark处理引擎进行处理。②窗口间隔窗口间隔指窗口的持续时间在窗口操作中只有窗口内的数据长度满足条件时才会触发批数据的处理。③滑动间隔滑动间隔Slide Duration指经过多长时间窗口滑动一次形成新的窗口这里必须注意的是滑动间隔和窗口间隔的大小一定要设置为批处理间隔的整数倍。批处理间隔、窗口间隔、滑动间隔的关系如图10-13所示。批处理间隔是1个时间单位窗口间隔是3个时间单位滑动间隔是2个时间单位。只有当窗口间隔满足条件时DStream才触发数据处理。常用的窗口转换操作如下。①windowwindowLengthslideInterval​返回一个基于原DStream的窗口批次计算后得到的新DStream。②countByWindowwindowLengthslideInterval​返回滑动窗口内DStream中元素的数量。③reduceByWindowfuncwindowLengthslideInterval​基于滑动窗口对原DStream中的元素进行聚合操作得到一个新DStream。④reduceByKeyAndWindowfuncwindowLengthslideInterval​[numTasks]​​基于滑动窗口键值对类型的DStream中的值按Key使用聚合函数func进行聚合操作得到一个新DStream。⑤reduceByKeyAndWindowfuncinvFuncwindowLengthslideInterval​[numTasks]​​一个更高效的reduceByKeyAndWindow​的实现版本先对滑动窗口中新的时间间隔内的数据增量聚合并去除最早的与新增数据量的时间间隔内的数据统计量。例如计算T4时刻过去5s窗口的WordCount那么可以将T3时刻过去5s的统计量加上[T3T4]的统计量再减去[T-2T-1]的统计量这种方法可以复用中间3s的统计量提高统计的效率。⑥countByValueAndWindowwindowLengthslideInterval​[numTasks]​​基于滑动窗口计算原DStream中每个RDD内每个元素出现的频次并返回DStream[​KeyLong​]​。其中Key是RDD中元素的类型Long是元素频次。与countByValue一样reduce任务的并发数可以通过一个可选参数进行配置。3输出操作。Spark Streaming在计算完成后可使用DStream将计算结果输出到外部系统如数据库或文件系统常用的方法如下。①print​​在Driver中打印出DStream中数据的前10个元素。②saveAsTextFilesprefix​[suffix]​​将DStream中的内容以文本的形式保存为文本文件其中每次批处理间隔内产生的文件都以prefix-TIME_IN_MS[.suffix]的方式命名。③saveAsObjectFilesprefix​[suffix]​​将DStream中的内容序列化并且以SequenceFile的格式保存其中每次批处理间隔内产生的文件都以prefix-TIME_IN_MS[.suffix]的方式命名。④saveAsHadoopFilesprefix​[suffix]​​将DStream中的内容以文本的形式保存为Hadoop文件其中每次批处理间隔内产生的文件都以prefix-TIME_IN_MS[.suffix]的方式命名。⑤foreachRDDfunc​将func函数应用于DStream的RDD上这个操作会把数据输出到外部系统比如保存RDD到文件或者数据库等。需要注意的是func函数是在运行该Streaming应用的Driver进程中被执行的。2.3.6、DStream的数据持久化与RDD一样DStream也能通过persist​将数据流缓存在内存中默认持久化方式是MEMORY_ONLY_SER也就是以序列化的方式将数据存放在内存中这样做的好处是遇到需要多次迭代计算的程序时可以直接使用内存中的数据速度优势十分明显。对于一些基于窗口的操作如reduceByWindow、reduceByKeyAndWindow以及基于状态的操作如updateStateBykey默认持久化策略为保存在内存中。对于来自外部的数据源Kafka、Flume、Sockets等​默认持久化策略是将数据副本保存在其他两台机器上。另外对于窗口和有状态的操作Spark应用程序必须配置CheckPoint通过StreamingContext来设置CheckPoint目录通过DStream设置CheckPoint间隔时间间隔必须是滑动间隔的倍数。2.3.7、Spark Streaming的性能优化Spark Streaming的性能优化主要包括优化运行时间和优化内存两方面。1优化运行时间。1增加并行度确保使用整个集群的资源而不是把任务集中在几个特定的节点上。对于包含Shuffle的操作增加其并行度以确保更充分地使用集群资源。2减少数据序列化、反序列化的负担SparkStreaming默认将接收的数据序列化后存储以减少内存的使用。但是序列化和反序列化需要更多CPU时间因此更加高效的序列化方式Kryo和自定义的序列化接口可以更高效地使用CPU。3设置合理的批处理间隔在Spark Streaming中Job之间存在依赖关系后面的Job必须确保前面的Job 执行结束后才能提交。若前面的Job执行的时间超出了批处理间隔那么后面的Job将无法按时提交这样就会进一步延迟接下来的Job执行造成后续Job的阻塞。因此需要设置一个合理的批处理间隔以确保Job能够在这个批处理间隔内结束批处理计算。4减少因任务提交和分发带来的负担通常情况下Akka框架能够高效地确保任务及时分发但是当批处理间隔非常小500ms时提交和分发任务的延迟就变得不可接受。使用Standalone和Coarse-grained Mesos模式通常会比使用Fine-grained Mesos模式有更小的延迟。2优化内存使用。1控制批处理间隔内的数据量Spark Streaming会把批处理间隔内接收的所有数据存放在Spark内部的可用内存区域中因此必须确保当前节点Spark的可用内存至少能容纳批处理间隔内的所有数据否则必须增加新的内存资源以提高集群的处理能力。2及时清理不再使用的数据Spark Streaming会将接收的数据全部存储到内存区域中因此对于处理过后就不再需要的数据应及时清理以确保Spark Streaming有更多的可用内存空间。通过设置合理的spark.cleaner.ttl时长来及时清理超时的无用内存数据这个参数需要小心设置以免后续操作中所需要的数据被当作超时数据清理掉。3调整GC策略GC会影响Job的正常运行延长Job的执行时间引起一系列不可预料的问题。观察GC的运行情况采用不同GC策略以进一步减小内存回收对Job运行的影响。2.4、Spark SQL、DataFrame、DataSet的使用2.4.1、Spark SQL简介Spark SQL是用于结构化数据处理的Spark模块。它提供了两个编程抽象分别叫作DataFrame和DataSet它们均用于分布式SQL查询引擎。在外部Spark SQL会将无论是DataFrame还是DataSet定义的数据都抽象为结构化数据这样使得开发人员可以像使用简单的SQL语句一样在Spark上轻松完成复杂的大数据交互操作。在内部Spark SQL将SQL执行操作转换成RDD然后提交到集群执行。2.4.2、Spark SQL查询语句Spark SQL最简单的用法是直接执行SQL查询语句可以使用最基本的SQL语法也可以选择HiveSQL语法。Spark SQL可以从已有的Hive中读取数据如果用其他编程语言运行SQL则Spark SQL将以DataFrame的形式返回结果。2.4.3、DataFrameDataFrame是一种分布式数据集合每一条数据都由多个字段组成。从概念上来说它与关系型数据库的表或者R语言和Python中的DataFrame等价只不过在底层DataFrame采取了更多优化。DataFrame可以从很多数据源加载数据如结构化数据文件、Hive表、关系型数据库、NoSQL数据库或者已有的RDD。2.4.4、DataSetDataSet的目的是把RDD的优势强类型可以使用lambda表达式函数处理数据和Spark SQL的优化执行引擎的优势结合到一起。DataSet由Java对象构建在DataSet上可以使用各种Transformation算子如map、flatMap、filter等来完成数据的交互式计算。2.4.5、创建一个Spark SQL应用在上述Spark项目的基础上创建Spark SQL项目具体步骤如下。1添加Maven依赖。打开pom.xml文件并将Spark SQL依赖加入项目具体代码如下。2新建SQLSimple类。在项目的Java目录下新建一个名为SQLSimple的Java类并在类中输入下面代码构建一个简单的Spark SQL程序。上述代码定义了一个简单的Spark SQL应用。具体步骤为定义SparkSession定义DataSet或DataFrame应用程序可以从数据文件、现有RDD、Hive表或Spark数据源创建上述代码通过JSON文件创建一个DataSet使用DataSet或DataFrame应用程序可以基于DataFrame做表结构的展示或者查询操作。JSON文件的数据具体如下。3打包和作业提交。在项目的根目录下输入如下命令对项目进行打包打包后的程序在Target目录下。Spark SQL的作业提交与一般的作业提交类似具体代码如下。在作业提交后通过日志可以看到打印出以下Spark SQLDataSet中的数据和Schema信息。4其他简单的操作。①查询name列上的数据具体代码和输出结果如下。②基于age列做groupBy分组然后执行count统计操作具体代码和输出结果如下。除了简单的字段引用和表达式支持DataFrame还提供了丰富的工具函数库包括字符串组装、日期处理、常见的数学函数等。2.4.6、Spark SQL的视图操作SparkSession上的SQL函数可以通过编程的方式运行SQL语句并将结果作为DataSet返回。应用程序可以使用如下方式将DataFrame转化为一个视图并基于视图进行查询操作。从使用层面上SparkSQL的视图可以理解为与数据库视图类似只不过数据库视图是对数据库表的抽象而Spark SQL的视图是对DataFrame的抽象。Spark SQL的视图可分为临时视图和全局临时视图。1临时视图Spark SQL的临时视图是会话范围的 如果创建它的SparkSession终止则它将消失。创建临时视图的代码如下。2全局临时视图如果希望拥有一个在所有会话之间共享的临时视图并保持活动状态直到Spark应用程序终止则可以通过创建全局临时视图的方式来实现。全局临时视图与系统保留的数据库global_temp绑定应用程序必须使用限定名称来引用它。具体使用如下。2.4.7、创建DataSetSaprk SQL DataSet API提供的API与RDD类似但是DataSet并没有使用Kryo实现序列化而是使用Spark提供的Encoder编码器序列化对象以便通过网络进行处理或传输。虽然Encoder和标准序列化都负责将对象转换为字节但Encoder是动态生成代码的并允许Spark执行更多操作例如过滤、排序和散列​而不需要将字节数组反序列化为对象。通过Encoder创建数据集的具体步骤如下。1创建Java Bean创建基于Java Bean的数据结构Spark会将Java Bean中的属性映射为视图的字段也可以理解为表的字段。下面定义一个简单的Persion类用于表示Java Bean。2创建Encoder基于Java Bean创建对应的Encoder。3创建DataSet并使用根据Encoder和Java Bean创建DataSet这个时候应用程序就可以基于DataSet做交互式查询操作。创建Encoder和DataSet的代码如下。除了基于Java Bean创建具有明确属性的数据结构应用程序还可以使用Encoder创建基本数据类型的结构具体实现如下。除了基于Java Bean创建DataSet应用程序还可以从JSON文件创建DataSet具体实现如下。2.4.8、DataSet与RDD相互转换1利用反射来推断SchemaSpark SQL支持两种不同方法将现有的RDD转换为DataSet。第一种方法利用反射机制来推断包含特定类型对象的RDD的Schema。这种基于反射的方法可以提供更简洁的代码这种模式要求编写Spark应用程序是在已经了解数据模型的前提下。具体实现代码如下。2以编程方式指定Schema如果无法提前定义JavaBean的数据类型则可以通过编程方式创建DatasetRow其创建过程分为以下3步。①从已有的RDD创建一个包含Row对象的RDD。如下代码创建一个名为peopleRDD1的RDD。其中textpath内的数据如下。②通过StructType创建一个Schema与步骤①中创建的RDD的结构相匹配。在如下代码中根据fields字段列表定义的Schema的结构与根据textpath构建的peopleRDD1相对应都包含name和age两个属性。③通过SparkSession提供的createDataFrame方法将Schema应用于rowRDD然后应用程序就可以基于该DataFrame执行SQL查询。具体实现代码如下。3DataFrames聚合操作除了基本的转换操作DataFrame函数还提供了一些聚合操作例如count​​、countDistinct​​、avg​​、max​​、min​等。此外除了使用Spark预定义的聚合函数应用程序可以创建自己的聚合函数。2.4.9、DataFrame数据的加载DataFrame支持以文本文件、JSON、Parquet、ORC等多种数据源的方式加载和保存数据。1读取和写入Parquet文件下面代码从Parquet文件加载数据查询出name和age对应的数据后再保存到新的Parquet文件中。2手动指定文件类型除了直接调用read​加载数据源应用程序还可以手动指定要加载的数据源。常用的数据源类型有JSON、Parquet、JDBC、ORC、LIBSVM、CSV、Text。如下代码加载JSON文件查询出name和age对应的数据后再保存到Parquet文件中。3直接在文件上运行SQL应用程序可以直接使用SQL查询该文件而不是通过读取API将文件加载到DataFrame再进行查询。如下代码直接基于Parquet文件执行SQL查询语句。2.4.10、DataFrame数据的保存在基于DataFrame完成计算后可以将计算结果保存到内存或外部存储系统。1数据保存模式DataFrame保存操作可以使用SaveMode设置保存数据的模式具体有如下保存模式。SaveMode.ErrorIfExistsdefault​默认模式当从DataFrame向数据源保存数据时如果数据已经存在则抛出异常。SaveMode.Append如果数据或表已经存在则将DataFrame的数据追加到已有数据的尾部。SaveMode.Overwrite如果数据或表已经存在则使用DataFrame数据覆盖之前的数据。SaveMode.Ignore如果数据已经存在则放弃保存DataFrame数据。这与SQL里的CREATE TABLE IF NOTEXISTS类似。2保存到持久化表在使用HiveContext的时候DataFrame可以用saveAsTable方法将数据保存成持久化表。与registerTempTable不同saveAsTable会将DataFrame的实际数据保存下来并且在HiveMetastore中创建一个游标指针。持久化表会一直保留即使Spark程序重启也不受影响只要程序连接到同一个Metastore就可以读取其数据。当读取持久化表时只需要把表名作为参数调用SQLContext.table方法即可得到对应的DataFrame。2.5、Spark Structured Streaming的使用2.5.1、Spark Structured Streaming简介Spark Structured StreamingSpark结构化流是基于Spark SQL引擎扩展的流处理引擎使得用户可以基于SQL像处理静态数据一样处理流式计算。Spark SQL引擎将负责不断地运行Structured Streaming并在流数据持续到达时更新最终结果。同时Spark StructuredStreaming通过检查点和预写日志确保端到端的一次性容错保证。简而言之Spark Structured Streaming提供快速、可扩展、容错、端到端的精确一次流处理而不需要用户关心具体的复杂流处理实现。在默认情况下Spark Structured Streaming使用微批处理引擎进行处理该引擎将数据流作为一系列小批量作业处理从而实现低至100ms的端到端延迟和完全一次的容错保证。从Spark 2.3后引入了一种被称为Continuous Processing的新的低延迟处理模式可以实现低至1ms的端到端延迟并且提供至少一次的容错保证。2.5.2、Spark Structured Streaming的数据模型Spark Structured Streaming的核心思想是将连续不断的数据流看作一个不断被连续追加的无界数据表这样源源不断的流式数据计算将被抽象为基于增量数据表的SQL操作。用户只需要定义SQL操作的方法Spark StructuredStreaming会在有增量数据时将增量数据收集起来组成微批数据并在该数据上执行SQL操作以完成基于SQL的流式计算。Spark Structured Streaming的数据模型如图所示。Spark Structured Streaming的数据流被看成表的行数据连续地向表中追加。Spark Structured Streaming查询将会产生一个结果表Result Table​第一行是时间线每秒有一个触发器第二行是输入流对输入流执行查询后产生的结果最终会被更新到第三行的结果表中。第四行是查询结果输出。Spark Structured Streaming的查询流程如图所示。Spark Structured Streaming的查询结果输出有3种不同的模式。完全模式Complete Mode​将计算结果更新整个结果表。追加模式Append Mode​将上一次触发到当前时间段内的数据追加到结果表的新行并写入外部存储。更新模式Update Mode​将上一次触发到当前时间段内的数据在结果表中更新的行写入外部存储。这种模式不同于完全模式它仅仅输出上一次触发后改变的行。2.5.3、创建一个Spark Structured Streaming应用在上述Spark项目的基础上创建Spark StructuredStreaming项目具体步骤如下。1新建SSSimple类。在Spark项目的Java目录下新建一个名为SSSimpleStructured Streaming Simple的Java类并在类中输入下面代码构建一个简单的Spark StructuredStreaming程序。该程序监听TCP端口上的数据并进行实时处理。上述代码定义了一个简单的Spark Structured Streaming应用。具体步骤为定义SparkSession定义DataFrame注意这里使用spark.readStream​监听TCP端口的数据并实时转换为DataSet定义数据操作这里基于DataFrame做聚合操作启动StreamingQuery实时流计算。在上述代码中名为lines的DataFrame表示一个包含流文本数据的无界表该表包含一列名为value的字符串并且流式数据中的每一条数据都将映射为表中的一行数据。然后使用lines.asEncoders.STRING​​将DataFrame转换为名为words的Dataset以便应用程序可以应用flatMap操作将每行都拆分为多个单词。words中包含了所有单词。接着通过words.groupBy“value”.count​对Dataset中的Value值进行分组统计操作并将计算结果定义成名为wordCounts的DataFrame。注意wordCounts是一个流式DataFrame它实现流式数据的实时查询。最后通过调用start​方法启动该流式计算。2启动TCP服务端。在Linux上通过Netcat启动一个TCP服务端具体代码如下。3打包和作业提交。在项目的根目录下输入如下命令对项目进行打包打包后的程序在Target目录下。SparkStreaming的作业提交与一般的作业提交类似具体代码如下。在Netcat窗口中输入“hello java”​查看Spark日志看到输出以下统计结果。通过上述日志我们不但能看到具体的执行结果还能明确地看到Spark内部Streaming Query的执行过程。该执行过程不但记录了runId和batchId等流式计算的描述信息同时定义了stateOperator列表用于状态监控定义了source表示实时查询的数据源定义了sink表示最终数据处理结果的输出。
返回列表