Advertisement

【Flink】Trigger操作

阅读量:

定义 了一个窗口,但是要求每来两条消息触发一次计算

复制代码
      .keyBy(_.id)
      .timeWindow(Time.days(1))
      .trigger(new TTrigger)
      .sum(1)
      .process(new BankProcess)
      .print()
    
    =============================自定义Trigger=========================================
     class TTrigger extends Trigger[Chenji, Window] {
    //定义一个状态用来保存消息的条数
    lazy val count: ValueStateDescriptor[Int] = new ValueStateDescriptor[Int]("count", classOf[Int])
    
    override def onElement(t: Chenji, l: Long, w: Window, triggerContext: Trigger.TriggerContext): TriggerResult = {
      val value: ValueState[Int] 

全部评论 (0)

还没有任何评论哟~