aboutsummaryrefslogtreecommitdiff
path: root/core/operator
diff options
context:
space:
mode:
Diffstat (limited to 'core/operator')
-rw-r--r--core/operator/keyby.go32
-rw-r--r--core/operator/map.go53
-rw-r--r--core/operator/operator.go16
-rw-r--r--core/operator/stateful_map.go75
4 files changed, 176 insertions, 0 deletions
diff --git a/core/operator/keyby.go b/core/operator/keyby.go
new file mode 100644
index 0000000..1296df0
--- /dev/null
+++ b/core/operator/keyby.go
@@ -0,0 +1,32 @@
+package operator
+
+import (
+ "context"
+ "goflink/core"
+)
+
+type KeySelector[IN any] func(in IN) string
+type KeyByOperator[IN any] struct {
+ selector KeySelector[IN]
+}
+
+func NewKeyByOperator[IN any](selector KeySelector[IN]) *KeyByOperator[IN] {
+ return &KeyByOperator[IN]{selector: selector}
+}
+
+func (k *KeyByOperator[IN]) Open(ctx context.Context) error {
+ return nil
+}
+
+func (k *KeyByOperator[IN]) Process(ctx context.Context, record core.StreamRecord[IN], out Emitter[IN]) error {
+ if record.Type != core.RecordData {
+ return out.Emit(ctx, record)
+ }
+
+ record.Key = k.selector(record.Value)
+ return out.Emit(ctx, record)
+}
+
+func (k *KeyByOperator[IN]) Close(ctx context.Context) error {
+ return nil
+}
diff --git a/core/operator/map.go b/core/operator/map.go
new file mode 100644
index 0000000..ef6912e
--- /dev/null
+++ b/core/operator/map.go
@@ -0,0 +1,53 @@
+package operator
+
+import (
+ "context"
+ "fmt"
+ "goflink/core"
+)
+
+// 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 core.StreamRecord[IN], out Emitter[OUT]) error {
+ if record.Type != core.RecordData {
+ // watermark and barrier
+ return out.Emit(ctx, core.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, core.StreamRecord[OUT]{
+ Type: core.RecordData,
+ Key: record.Key,
+ Timestamp: record.Timestamp,
+ Value: outValue,
+ })
+}
+
+func (fo *MapOperator[IN, OUT]) Close(ctx context.Context) error {
+ return nil
+}
diff --git a/core/operator/operator.go b/core/operator/operator.go
new file mode 100644
index 0000000..d4bf21b
--- /dev/null
+++ b/core/operator/operator.go
@@ -0,0 +1,16 @@
+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
+}
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)
+}