diff options
| author | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-09 18:32:44 +0700 |
|---|---|---|
| committer | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-09 18:32:44 +0700 |
| commit | 4ca27aec303f54dc9ee670c67c2fefa80dd004a1 (patch) | |
| tree | 5628714f47e374bfb98c1cea99970b19892e8b3e /core/operator/stateful_map.go | |
| parent | a892cb5b1d121177881fa56a79abbafef5880a80 (diff) | |
| download | goflink-4ca27aec303f54dc9ee670c67c2fefa80dd004a1.tar.gz goflink-4ca27aec303f54dc9ee670c67c2fefa80dd004a1.zip | |
feat: add KeyByOperator and wordCounter for stateful processing
Diffstat (limited to 'core/operator/stateful_map.go')
| -rw-r--r-- | core/operator/stateful_map.go | 75 |
1 files changed, 75 insertions, 0 deletions
diff --git a/core/operator/stateful_map.go b/core/operator/stateful_map.go new file mode 100644 index 0000000..7263170 --- /dev/null +++ b/core/operator/stateful_map.go @@ -0,0 +1,75 @@ +package operator + +import ( + "context" + "fmt" + + "goflink/core" +) + +// RichMapFunc is a stateful map on the user side: Open grabs the state handles +// off ctx, Map runs per record with the backend already pointed at that +// record's key. +type RichMapFunc[IN, OUT any] interface { + Open(ctx context.Context) error + Map(ctx context.Context, in IN) (OUT, error) + Close(ctx context.Context) error +} + +type StatefulMapOperator[IN, OUT any] struct { + userFunc RichMapFunc[IN, OUT] + backend core.StateBackend +} + +func NewStatefulMapOperator[IN, OUT any](fn RichMapFunc[IN, OUT]) *StatefulMapOperator[IN, OUT] { + return &StatefulMapOperator[IN, OUT]{userFunc: fn} +} + +func (s *StatefulMapOperator[IN, OUT]) Open(ctx context.Context) error { + // Get backend state from context via using context KV + be, ok := core.ExtractStateBackend(ctx) + if !ok || be == nil { + return fmt.Errorf("stateful map: no state backend in context") + } + + s.backend = be + // The user func opens its ValueState off the same ctx. + return s.userFunc.Open(ctx) +} + +func (s *StatefulMapOperator[IN, OUT]) Process(ctx context.Context, record core.StreamRecord[IN], out Emitter[OUT]) error { + if record.Type != core.RecordData { + return out.Emit(ctx, core.StreamRecord[OUT]{ + Type: record.Type, + Timestamp: record.Timestamp, + Key: record.Key, + BarrierID: record.BarrierID, + }) + } + + // Without a key every record would silently share the "" state slot. + if record.Key == "" { + return fmt.Errorf("stateful map: record %v has no key, keyBy it first", record.Value) + } + + if err := s.backend.SetCurrentKey(record.Key); err != nil { + return fmt.Errorf("set current key %q: %w", record.Key, err) + } + + // run the user transform + res, err := s.userFunc.Map(ctx, record.Value) + if err != nil { + return fmt.Errorf("cannot process value %v: %w", record.Value, err) + } + + return out.Emit(ctx, core.StreamRecord[OUT]{ + Type: core.RecordData, + Timestamp: record.Timestamp, + Key: record.Key, + Value: res, + }) +} + +func (s *StatefulMapOperator[IN, OUT]) Close(ctx context.Context) error { + return s.userFunc.Close(ctx) +} |