【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)
还没有任何评论哟~
