package operator import ( "context" "goflink/core" ) // TimestampAssigner is responsible for event time defined by user type TimestampAssigner[IN any] func(in IN) int64 type WaterGeneratorOperator[IN any] struct { assigner TimestampAssigner[IN] maxOutOfOrder int64 currentMaxTimestamp int64 lastEmittedWatermark int64 } func NewWatermarkGeneratorOperator[IN any](maxOutOfOrder int64, assigner TimestampAssigner[IN]) *WaterGeneratorOperator[IN] { return &WaterGeneratorOperator[IN]{ maxOutOfOrder: maxOutOfOrder, assigner: assigner, lastEmittedWatermark: -1, } } func (w *WaterGeneratorOperator[IN]) Open(ctx context.Context) error { return nil } func (w *WaterGeneratorOperator[IN]) Process(ctx context.Context, record core.StreamRecord[IN], out Emitter[IN]) error { if record.Type == core.RecordData { // query the time eventTime := w.assigner(record.Value) record.Timestamp = eventTime // update the max timestamp if w.currentMaxTimestamp < eventTime { w.currentMaxTimestamp = eventTime } // calc the watermark watermark := w.currentMaxTimestamp - w.maxOutOfOrder // check if current watermark later than previous watermark if watermark > w.lastEmittedWatermark { w.lastEmittedWatermark = watermark return out.Emit(ctx, core.NewWatermarkRecord[IN](watermark)) } return nil } return out.Emit(ctx, record) } func (w *WaterGeneratorOperator[IN]) Close(ctx context.Context) error { return nil }