最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java 數(shù)據(jù)流之Broadcast State

 更新時(shí)間:2021年09月14日 10:44:56   作者:Vicky_Tang  
這篇文章主要介紹了Java 數(shù)據(jù)流之Broadcast State,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下

一、BroadcastState 的介紹

廣播狀態(tài)(Broadcast State)是 Operator State 的一種特殊類型。如果我們需要將配置 、規(guī)則等低吞吐事件流廣播到下游所有 Task 時(shí),就可以使用 BroadcastState。下游的 Task 接收這些配置、規(guī)則并保存為 BroadcastState,所有Task 中的狀態(tài)保持一致,作用于另一個(gè)數(shù)據(jù)流的計(jì)算中。
簡(jiǎn)單理解:一個(gè)低吞吐量流包含一組規(guī)則,我們想對(duì)來自另一個(gè)流的所有元素基于此規(guī)則進(jìn)行評(píng)估。
場(chǎng)景:動(dòng)態(tài)更新計(jì)算規(guī)則。

廣播狀態(tài)與其他操作符狀態(tài)的區(qū)別在于:

  • 它有一個(gè) map 格式,用于定義存儲(chǔ)結(jié)構(gòu)
  • 它僅對(duì)具有廣播流和非廣播流輸入的特定操作符可用
  • 這樣的操作符可以具有不同名稱的多個(gè)廣播狀態(tài)

二、BroadcastState 操作流程

三、案例實(shí)現(xiàn)

  • 從端口讀取Json數(shù)據(jù)作為事件流
  • 從Mysql讀取數(shù)據(jù)作為廣播流
  • 關(guān)聯(lián)廣播流和事件流
  • 匹配對(duì)應(yīng)的用戶信息
package cn.kgc.broadcast
 
import java.sql.{Connection, DriverManager, PreparedStatement}
 
import com.alibaba.fastjson.JSON
import org.apache.flink.api.common.state.{BroadcastState, MapStateDescriptor}
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.datastream.BroadcastStream
import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction
import org.apache.flink.streaming.api.functions.source.{RichParallelSourceFunction, SourceFunction}
import org.apache.flink.streaming.api.scala._
import org.apache.flink.util.Collector
 
// (001,'tom',18,'北京',15830010002)
// 定義樣例類 接受 MySQL的用戶數(shù)據(jù)
case class BaseUserInfo(id:Long,name:String,age:Int,city:String,phone:Long)
 
// user_id、user_name、user_addrss、behaviour、url
// 輸出數(shù)據(jù)類型
case class UserVisitInfo(id:Long,name:String,city:String,behaviour:String,url:String)
 
// 實(shí)現(xiàn)廣播ProcessFunction
class MyBroadcastFunc extends BroadcastProcessFunction[String,(Long, BaseUserInfo),UserVisitInfo]{
 
  lazy val mapStateDes = new MapStateDescriptor[Long, BaseUserInfo]("mapState",classOf[Long],classOf[BaseUserInfo])
 
  // 處理的是日志流中的每條數(shù)據(jù)
  override def processElement(value: String, ctx: BroadcastProcessFunction[String, (Long, BaseUserInfo), UserVisitInfo]#ReadOnlyContext, out: Collector[UserVisitInfo]): Unit = {
    // {"user_id":"001","ts":"2021-07-10 11:10:05","behaviour":"browse","url":"https://www.tb1.com/1.html"}
    val user_id = JSON.parseObject(value).getLong("user_id")
    val behaviour = JSON.parseObject(value).getString("behaviour")
    val url = JSON.parseObject(value).getString("url")
 
    val mapState = ctx.getBroadcastState(mapStateDes)
    val userInfo = mapState.get(user_id)
 
    out.collect(UserVisitInfo(user_id,userInfo.name,userInfo.city,behaviour,url))
 
  }
 
  // 處理的是廣播流的每個(gè)值
  override def processBroadcastElement(value: (Long, BaseUserInfo), ctx: BroadcastProcessFunction[String, (Long, BaseUserInfo), UserVisitInfo]#Context, out: Collector[UserVisitInfo]): Unit = {
    val mapState: BroadcastState[Long, BaseUserInfo] = ctx.getBroadcastState(mapStateDes)
    mapState.put(value._1,value._2)
  }
}
 
 
class UserSourceFunc extends RichParallelSourceFunction[BaseUserInfo]{
 
  var conn:Connection = _
  var statement: PreparedStatement = _
  var flag:Boolean = true
 
  override def open(parameters: Configuration): Unit = {
    conn = DriverManager.getConnection("jdbc:mysql://localhost:3306/test?characterEncoding=utf-8&serverTimezone=UTC","root","liu911223")
    statement = conn.prepareStatement("select * from base_user")
  }
 
  override def run(ctx: SourceFunction.SourceContext[BaseUserInfo]): Unit = {
    while (flag){
      Thread.sleep(5000)
      val resultSet = statement.executeQuery()
      while (resultSet.next()){
        val id = resultSet.getLong(1)
        val name = resultSet.getString(2)
        val age = resultSet.getInt(3)
        val city = resultSet.getString(4)
        val phone = resultSet.getLong(5)
        ctx.collect(BaseUserInfo(id,name,age,city,phone))
      }
    }
  }
 
  override def cancel(): Unit = {
    flag = false
  }
 
  override def close(): Unit = {
    if (statement != null) statement.close()
    if (conn != null) conn.close()
  }
}
object BroadcastDemo01 {
  def main(args: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setParallelism(1)
 
    // 定義為KV,一方面是為了廣播的時(shí)候定義為map,另一方面是為了做關(guān)聯(lián)操作
    val userBaseDS: DataStream[(Long, BaseUserInfo)] = env.addSource(new UserSourceFunc)
      .map(user => (user.id, user))
    val mapStateDes = new MapStateDescriptor[Long, BaseUserInfo]("mapState",classOf[Long],classOf[BaseUserInfo])
    val broadCastStream: BroadcastStream[(Long, BaseUserInfo)] = userBaseDS.broadcast(mapStateDes)
 
    // 日志JSON數(shù)據(jù)
    val dataInfoDS: DataStream[String] = env.socketTextStream("master",1314)
 
    dataInfoDS.connect(broadCastStream)
      .process(new MyBroadcastFunc)
      .print()
 
    env.execute()
  }
}

到此這篇關(guān)于Java 數(shù)據(jù)流之Broadcast State的文章就介紹到這了,更多相關(guān)Java Broadcast State內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java中數(shù)組與集合的相互轉(zhuǎn)換實(shí)現(xiàn)解析

    Java中數(shù)組與集合的相互轉(zhuǎn)換實(shí)現(xiàn)解析

    這篇文章主要介紹了Java中數(shù)組與集合的相互轉(zhuǎn)換實(shí)現(xiàn)解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-08-08
  • RabbitMQ消費(fèi)者限流實(shí)現(xiàn)消息處理優(yōu)化

    RabbitMQ消費(fèi)者限流實(shí)現(xiàn)消息處理優(yōu)化

    這篇文章主要介紹了RabbitMQ消費(fèi)者限流實(shí)現(xiàn)消息處理優(yōu)化,消費(fèi)者限流是用于消費(fèi)者每次獲取消息時(shí)限制條數(shù),注意前提是手動(dòng)確認(rèn)模式,并且在手動(dòng)確認(rèn)后才能獲取到消息,感興趣想要詳細(xì)了解可以參考下文
    2023-05-05
  • swagger2隱藏在API文檔顯示某些參數(shù)的操作

    swagger2隱藏在API文檔顯示某些參數(shù)的操作

    這篇文章主要介紹了swagger2隱藏在API文檔顯示某些參數(shù)的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • java方法重載示例

    java方法重載示例

    方法重載是以統(tǒng)一的方式處理不同數(shù)據(jù)類型的一種手段,這篇文章主要介紹了java方法重載示例,需要的朋友可以參考下
    2014-03-03
  • Spring中@Autowired @Resource @Inject三個(gè)注解有什么區(qū)別

    Spring中@Autowired @Resource @Inject三個(gè)注解有什么區(qū)別

    在我們使用Spring框架進(jìn)行日常開發(fā)過程中,經(jīng)常會(huì)使用@Autowired, @Resource, @Inject注解來進(jìn)行依賴注入,下面來介紹一下這三個(gè)注解有什么區(qū)別
    2023-03-03
  • SpringSecurity 認(rèn)證、注銷、權(quán)限控制功能(注銷、記住密碼、自定義登入頁)

    SpringSecurity 認(rèn)證、注銷、權(quán)限控制功能(注銷、記住密碼、自定義登入頁)

    SpringSecurity是一個(gè)強(qiáng)大的Java框架,用于保護(hù)應(yīng)用程序的安全性,它提供了一套全面的安全解決方案,本文給大家介紹SpringSecurity認(rèn)證、注銷、權(quán)限控制和注銷、記住密碼、自定義登入頁等知識(shí)總結(jié),感興趣的朋友一起看看吧
    2025-03-03
  • java中Object類4種方法詳細(xì)介紹

    java中Object類4種方法詳細(xì)介紹

    大家好,本篇文章主要講的是java中Object類4種方法詳細(xì)介紹,感興趣的同學(xué)趕快來看一看吧,對(duì)你有幫助的話記得收藏一下,方便下次瀏覽
    2022-01-01
  • java隨機(jī)字符補(bǔ)充版

    java隨機(jī)字符補(bǔ)充版

    今天在zuidaimai看到一個(gè)java隨機(jī)字符生成demo,正好要用,但發(fā)現(xiàn)不完整,重新整理一下,分享給有需要的朋友
    2014-01-01
  • Mybatis 中如何判斷集合的size

    Mybatis 中如何判斷集合的size

    這篇文章主要介紹了在Mybatis中判斷集合的size操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧
    2021-02-02
  • 基于Java Callable接口實(shí)現(xiàn)線程代碼實(shí)例

    基于Java Callable接口實(shí)現(xiàn)線程代碼實(shí)例

    這篇文章主要介紹了基于Java Callable接口實(shí)現(xiàn)線程代碼實(shí)例,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-08-08

最新評(píng)論

隆子县| 涟水县| 图木舒克市| 宁海县| 将乐县| 东乡族自治县| 伊通| 昌乐县| 云南省| 长阳| 渭源县| 昌平区| 鄄城县| 呼图壁县| 长治市| 准格尔旗| 图片| 民乐县| 邵阳县| 洪泽县| 富源县| 金乡县| 神池县| 湟中县| 若尔盖县| 定南县| 石阡县| 荔浦县| 且末县| 奉化市| 江西省| 平湖市| 海口市| 灵武市| 舞阳县| 通海县| 旬邑县| 呼伦贝尔市| 石家庄市| 于都县| 重庆市|