package operator import ( "context" "goflink/core" ) type Emitter[OUT any] interface { Emit(ctx context.Context, record core.StreamRecord[OUT]) error } type Operator[IN, OUT any] interface { Open(ctx context.Context) error Process(ctx context.Context, record core.StreamRecord[IN], out Emitter[OUT]) error Close(ctx context.Context) error }