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)
}
|