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 }