SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

发布时间:2021-12-07 11:05:20 作者:柒染
来源:亿速云 阅读:123

本篇文章给大家分享的是有关SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决,小编觉得挺实用的,因此分享给大家学习,希望大家阅读完这篇文章后可以有所收获,话不多说,跟着小编一起来看看吧。

前言

当我在测试SparkStreaming的状态操作mapWithState算子时,当我们设置timeout(3s)的时候,3s过后数据还是不会过期,不对此key进行操作,等到30s左右才会清除过期的数据。

百度了很久,关于timeout的资料很少,更没有解决这个问题的文章,所以说,百度也不是万能的,有时候还是需要靠自己。

所以我就在周末研究了一下,然后将结果整理了出来,希望能帮助大家更全面的理解Spark状态计算。

mapWithState

按理说Spark Streaming实时处理,数据就像流水,每个批次之间的数据都是独立的,处理完就处理完了,不留下任何状态。但是免不了一些有状态的操作,例如统计从流启动到现在,某个单词出现了多少次,所以状态操作就出现了。

状态操作分为updateStateByKey和mapWithState,两者有着很大的区别。简单的来说,前者每次输出的都是全量状态,后者输出的是增量状态。

过期原理

SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

过期这一块估计很多人开始都理解错了,我刚开始理解就是数据从出现,经过多少秒之后就会过期。其实不是,这里的过期指的是空闲时间。

注释大概是这个意思:timeout()传入一个时间间隔参数,如果一个key在大于此间隔没有此key的数据流入,则被认为是空闲的,就会单独调用一次mapWithState中的func来清除这些空闲数据状态。

先写结论

使用了timeout()之后,需要使用以下代码来在间隔内清除失效key。

stream.checkpoint(Seconds(6))

checkpoint的时候,会开启全面扫描,才会对state中的失效key进行清理。

测试

val conf = new SparkConf().setMaster("local[2]").setAppName("state")val ssc = new StreamingContext(conf, Seconds(3))ssc.checkpoint("./tmp")val streams: DStream[(String, Int)] = ssc.socketTextStream("localhost", 9999)  .map(x => (x, 1))val result = streams.mapWithState(StateSpec.function((k: String, v: Option[Int], state: State[Int]) => {   val count = state.getOption().getOrElse(0)println(k)println(v)var sum = 0if (!state.isTimingOut()) {     sum = count + v.get
          state.update(sum)} else {     println("timeout")}Option(sum)  })  .timeout(Seconds(3)))// 这行代码是触发清除机制的关键// result.checkpoint(Seconds(6))result.print()ssc.start()ssc.awaitTermination()

使用上面的代码进行测试,设置过期时间为3s。但是3s过后发现key并没有过期,也不会被清除,大概30S之后被清除。

在9999端口输入一个tom后,不再进行任何操作。测试结果如下:

tom
Some(1)
-------------------------------------------
Time: 1618228587000 ms
-------------------------------------------
Some(1)


tom
None
timeout
-------------------------------------------
Time: 1618228614000 ms
-------------------------------------------
Some(0)

从测试结果可以看出,从输入到清除大概是27s。

我们现在将注释的代码放开,每6s进行checkpoint一次,输入tom:

tom
Some(1)
-------------------------------------------
Time: 1618228497000 ms
-------------------------------------------
Some(1)

tom
None
timeout
-------------------------------------------
Time: 1618228506000 ms
-------------------------------------------
Some(0)

从生成到清除用了9秒,正好是过期时间 + 下一个窗口时间,触发了checkpoint。

猜想

第一次学状态操作的时候,就考虑如何去掉一些过期的key,通过timeout()的方法没有完成自己想法,从网上也没有找到解决方案,所以就暂且搁置在一边了。后来又回过头来考虑这个问题,然后根据自己的想法去猜想、去验证。

1. 我先看的是mapWithState()的返回值

SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

2. MapWithStateDStreamImpl

SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

每个Dstream的计算逻辑都在compute()中,这里是调用了internalStream的getOrCompute(),根据继承关系,调用的是父类Dstream的此方法:

SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

getOrCompute()主要功能为:计算、缓存、checkpoint。这里只需要记住几个地方:checkpointDuration,即checkpoint间隔,和调用了checkpoint()。其实真正的计算还是调用了compute(),接着去看compute()

3. InternalMapWithStateDStream

SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

compute()里面也调用了getOrCompute()方法,其实和上面调用的一样,都是Dstream的,这里主要看的是使用createFromRDD()生成的StateRDD。

4. MapWithStateRDD

这个StateRDD就是参与状态计算的数据集合,首先看它是如何生成的:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

再看看StateRDD的compute()是如何计算的:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

从compute()看出,当doFullScan为true的时候,才会触发过期key的清除,updateRecordWithData()负责全面扫描清除过期key

这不,思路就来了,我们只要找到开启FullScan的方法,不就可以自行触发清除机制了吗!

那么,我们先看看doFullScan的默认值:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

默认是没开启的,接着通过快捷键看看哪些地方使用了doFullScan:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

从图中看出,有两处代码修改了doFullScan,我们找到这两处代码:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

第一个基本上排除,那么就剩下第二个:checkpoint(),我们要知道的是,状态操作必须要checkpoint

还记得在2中的getOrCompute()吗,当checkpointDuration不为null的时候,调用checkpoint()。
我们来看3中InternalMapWithStateDStream是如何定义这个duration的:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

如图,sideDuration是窗口时间,乘以系数10就是默认的checkpoint时长,所以当我设置窗口为3s时,checkpoint周期就是30s,30s才会清理一次过期key。

而通过checkpoint(interval)可以设置checkpoint的间隔,所以覆盖了上面程序中默认的30s。
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

5.MapWithStateRDDRecord

最后提一提,FullScan是在这个类中开启的,所以先看看这个Record的注释介绍:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

意思就是负责存储StateRDD的状态KV,updateRecordWithData()负责清除过期的Record,我们来看看这个方法的实现:
SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决

removeTimedoutData就是是否开启全面扫描,即doFullScan的值。

以上就是SparkStreaming使用mapWithState时设置timeout()无法生效问题该怎么解决,小编相信有部分知识点可能是我们日常工作会见到或用到的。希望你能通过这篇文章学到更多知识。更多详情敬请关注亿速云行业资讯频道。

推荐阅读:
  1. SparkStreaming的实现和使用方法
  2. 大数据开发中Spark Streaming处理数据及写入Kafka

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

sparkstreaming mapwithstate timeout

上一篇:JavaScript如何实现计算两数的乘积功能

下一篇:Hyperledger fabric Chaincode开发的示例分析

相关阅读

您好,登录后才能下订单哦!

密码登录
登录注册
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》