大数据系列-SPARK-STREAMING流数据state

package com.test

import org.apache.spark.SparkConf
import org.apache.spark.streaming.dstream.{DStream, ReceiverInputDStream}
import org.apache.spark.streaming.{Seconds, StreamingContext}

//有状态state函数updateStateByKey
object SparkStreamingState {

  def main(args: Array[String]): Unit = {

    val sparkConf = new SparkConf().setAppName("SparkStreamingState").setMaster("local[*]")
    val streamingContext = new StreamingContext(sparkConf, Seconds(5))
    streamingContext.checkpoint("data/cpDir")

    val dstream: ReceiverInputDStream[String] = streamingContext.socketTextStream("localhost", 8600)
    val wordToMap: DStream[(String, Int)] = dstream.map((_, 1))
    val wordState: DStream[(String, Int)] = wordToMap.updateStateByKey((seq: Seq[Int]/*相同KEY的VALUE值*/, option: Option[Int]/*缓冲区中相同KEY的VALUE值*/) => {
      Option(seq.sum + option.getOrElse(0))
    })
    wordState.print()

    streamingContext.start()
    streamingContext.awaitTermination()

  }

}
Logo

GitCode 天启AI是一款由 GitCode 团队打造的智能助手,基于先进的LLM(大语言模型)与多智能体 Agent 技术构建,致力于为用户提供高效、智能、多模态的创作与开发支持。它不仅支持自然语言对话,还具备处理文件、生成 PPT、撰写分析报告、开发 Web 应用等多项能力,真正做到“一句话,让 Al帮你完成复杂任务”。

更多推荐