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/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/map.go')
| -rw-r--r-- | core/map.go | 52 |
1 files changed, 0 insertions, 52 deletions
diff --git a/core/map.go b/core/map.go deleted file mode 100644 index bbb8dcd..0000000 --- a/core/map.go +++ /dev/null @@ -1,52 +0,0 @@ -package core - -import ( - "context" - "fmt" -) - -// MapFunc is business logic for user implement -type MapFunc[IN, OUT any] func(in IN) (OUT, error) - -type MapOperator[IN, OUT any] struct { - userFunc MapFunc[IN, OUT] -} - -func NewMapOperator[IN, OUT any](fn MapFunc[IN, OUT]) *MapOperator[IN, OUT] { - return &MapOperator[IN, OUT]{ - userFunc: fn, - } -} - -func (fo *MapOperator[IN, OUT]) Open(ctx context.Context) error { - return nil -} - -func (fo *MapOperator[IN, OUT]) Process(ctx context.Context, record StreamRecord[IN], out Emitter[OUT]) error { - if record.Type != RecordData { - // watermark and barrier - return out.Emit(ctx, StreamRecord[OUT]{ - Type: record.Type, - Key: record.Key, - Timestamp: record.Timestamp, - BarrierID: record.BarrierID, - }) - } - - outValue, err := fo.userFunc(record.Value) - - if err != nil { - return fmt.Errorf("cannot process value %v: %w", record.Value, err) - } - - return out.Emit(ctx, StreamRecord[OUT]{ - Type: RecordData, - Key: record.Key, - Timestamp: record.Timestamp, - Value: outValue, - }) -} - -func (fo *MapOperator[IN, OUT]) Close(ctx context.Context) error { - return nil -} |