Get the App
SLTechnology News&Howtos  ›  Internet Technology  › 

How to realize the calculation of stateful stateful in flink

Shulou Source: shulou.com Published: 2022-06-01 06:12:59 09月29日 Update

Editor to share with you how to achieve stateful stateful calculation in flink, I believe most people do not know much about it, so share this article for your reference, I hope you can learn a lot after reading this article, let's go to know it!

Import org.apache.flink.api.common.functions.RichFlatMapFunctionimport org.apache.flink.api.common.state.ValueStateimport org.apache.flink.util.Collectorimport org.apache.flink.configuration.Configurationimport org.apache.flink.api.common.state.ValueStateDescriptorimport org.apache.flink.streaming.api.scala.StreamExecutionEnvironment class CountWindowAverage extends RichFlatMapFunction [(Long, Double), (Long, Double)] {private var sum: ValueState [(Long, Double)] = _ override def flatMap (input: (Long, Double), out: Collector [(Long) Double)]): Unit = {/ / access the state valueval tmpCurrentSum = sum.value / / If it hasn't been used before, it will be null val currentSum = if (tmpCurrentSum! = null) {tmpCurrentSum} else {(0L, 0d)} / / update the count val newSum = (currentSum._1 + 1, currentSum._2 + input._2) / / update the state sum.update (newSum) / / if the count reaches 2 Emit the average and clear the state if (newSum._1 > = 2) {out.collect ((input._1, newSum._2 / newSum._1)) / / clear the state / / sum.clear ()}} override def open (parameters: Configuration): Unit = {sum = getRuntimeContext.getState (new ValueStateDescriptor [(Long, Double)] ("average", classOf [(Long) Double)}} import org.apache.flink.streaming.api.scala.StreamExecutionEnvironmentimport org.apache.flink.api.scala._object ECountWindowAverage {def main (args: Array [String]): Unit = {val env = StreamExecutionEnvironment.getExecutionEnvironment env.fromCollection (List ((1L, 3D), (1L, 5d), (1L, 7d), (1L, 4d), (1L, 4d), (1L KeyBy (_. _ 1) .flatMap (new CountWindowAverage ()) .print () / * .keyBy (_. _ 1) .flatMap (new CountWindowAverage ()) .print () * / the printed output will be (1line 4) and (1jue 5) env.execute ("ExampleManagedState")} above are all the contents of the article "how to implement stateful stateful calculation in flink" Thank you for reading! I believe we all have a certain understanding, hope to share the content to help you, if you want to learn more knowledge, welcome to follow the industry information channel!

Tags: Status articles content not much most more knowledge industry information information channels channels references learning help Apple Docker Huawei Linux macOS MariaDB Microsoft MySQL NVidia OPPO Reno Shulou Information Docker Microsoft Shulou Tech Info vpn