aboutsummaryrefslogtreecommitdiff
path: root/core/operator/keyed_process.go
blob: b507d87b6d67d8e89aace08c7858c364bf0ff683 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
package operator

import (
	"context"
	"fmt"
	"goflink/core"
)

type KeyedProcesser[IN, OUT any] interface {
	Open(ctx context.Context) error
	ProcessElement(ctx context.Context, record core.StreamRecord[IN], out Emitter[OUT]) error
	Close(ctx context.Context) error
	OnTimer(ctx context.Context, timestamp int64, out Emitter[OUT]) error
}

type KeyProcessOperator[IN, OUT any] struct {
	userFunc    KeyedProcesser[IN, OUT]
	timeService core.InternalTimerService 
	backend     core.StateBackend
}

func NewKeyProcessOperator[IN, OUT any](userFunc KeyedProcesser[IN, OUT]) *KeyProcessOperator[IN, OUT] {
	return &KeyProcessOperator[IN, OUT]{
		userFunc:    userFunc,
		timeService: core.NewTimerService(),
	}
}

func (k *KeyProcessOperator[IN, OUT]) Open(ctx context.Context) error {
	// init backend
	be, ok := core.ExtractStateBackend(ctx)
	if !ok {
		return fmt.Errorf("state backend not found")
	}

	k.backend = be
	// inject the timer to context for user use
	ctx = core.InjectTimerService(ctx, k.timeService)
	return k.userFunc.Open(ctx)

}

func (k *KeyProcessOperator[IN, OUT]) Process(ctx context.Context, record core.StreamRecord[IN], out Emitter[OUT]) error {

	switch record.Type {
	case core.RecordData:
		if err := k.backend.SetCurrentKey(record.Key); err != nil {
			return err
		}

		return k.userFunc.ProcessElement(ctx, record, out)
	case core.RecordWatermark:
		// advance the timer
		triggerTimerQueue := k.timeService.AdvanceWatermark(record.Timestamp)

		for _, t := range triggerTimerQueue {

			if err := k.backend.SetCurrentKey(t.Key); err != nil {
				return err
			}
			if err := k.userFunc.OnTimer(ctx, t.Timestamp, out); err != nil {
				return err
			}
		}
		return out.Emit(ctx, core.StreamRecord[OUT]{
			Type:      core.RecordWatermark,
			Timestamp: record.Timestamp,
		})
	default:
		return out.Emit(ctx, core.StreamRecord[OUT]{
			Timestamp: record.Timestamp,
			BarrierID: record.BarrierID,
			Key:       record.Key,
		})
	}

}

func (k *KeyProcessOperator[IN, OUT]) Close(ctx context.Context) error {
	return k.userFunc.Close(ctx)
}