diff options
Diffstat (limited to 'core/operator/watermark.go')
| -rw-r--r-- | core/operator/watermark.go | 58 |
1 files changed, 58 insertions, 0 deletions
diff --git a/core/operator/watermark.go b/core/operator/watermark.go new file mode 100644 index 0000000..75390a8 --- /dev/null +++ b/core/operator/watermark.go @@ -0,0 +1,58 @@ +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 +} |