基于交易量的Apache Flink自定义窗口聚合技术问询
基于交易量聚合交易流的自定义Trigger实现
我想要把交易流聚合为以累计交易量为划分标准的窗口——这里的交易量指的是区间内所有交易的交易规模(size字段)之和。
我已经编写了一个自定义Trigger来实现这个窗口划分逻辑,核心代码片段如下:
case class Trade(key: Int, millis: Long, time: LocalDateTime, price: Double, size: Int) class VolumeTrigger(triggerVolume: Int, config: ExecutionConfig) extends Trigger[Trade, Window] { val LOG: Logger = LoggerFactory.getLogger(classOf[VolumeTrigger]) // 后续的Trigger逻辑实现代码... }
这个VolumeTrigger的核心思路是:每当窗口内累计的交易size总和达到预先设定的triggerVolume阈值时,就触发窗口的计算与结果输出,以此实现按交易量维度来划分聚合窗口的需求。
内容的提问来源于stack exchange,提问作者Daniel
相关产品推荐
相关产品推荐

