From acc0b6fbaace61a3befd7f726f9d80fb69391126 Mon Sep 17 00:00:00 2001 From: SangTran-127 Date: Sat, 8 Aug 2026 12:54:43 +0700 Subject: implement core transport and processing operators with receiver, emitter, filter, and map functionalities --- core/map.go | 49 +++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 49 insertions(+) create mode 100644 core/map.go (limited to 'core/map.go') diff --git a/core/map.go b/core/map.go new file mode 100644 index 0000000..c011e8d --- /dev/null +++ b/core/map.go @@ -0,0 +1,49 @@ +package core + +import "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 NewMapFuncOperator[IN, OUT any](fn MapFunc[IN, OUT]) *MapOperator[IN, OUT] { + return &MapOperator[IN, OUT]{ + userFunc: fn, + } +} + +func (fo *MapOperator[IN, OUT]) Open() error { + return nil +} + +func (fo *MapOperator[IN, OUT]) Process(record StreamRecord[IN], out Emitter[OUT]) error { + if record.Type != RecordData { + // watermark and barrier + return out.Emit(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", record.Value) + } + + return out.Emit(StreamRecord[OUT]{ + Type: RecordData, + Key: record.Key, + Timestamp: record.Timestamp, + Value: outValue, + }) +} + +func (fo *MapOperator[IN, OUT]) Close() error { + return nil +} -- cgit v1.2.3