aboutsummaryrefslogtreecommitdiff
path: root/core/operator/keyby.go
diff options
context:
space:
mode:
authorSangTran-127 <tranquangsang12.7@gmail.com>2026-08-09 18:32:44 +0700
committerSangTran-127 <tranquangsang12.7@gmail.com>2026-08-09 18:32:44 +0700
commit4ca27aec303f54dc9ee670c67c2fefa80dd004a1 (patch)
tree5628714f47e374bfb98c1cea99970b19892e8b3e /core/operator/keyby.go
parenta892cb5b1d121177881fa56a79abbafef5880a80 (diff)
downloadgoflink-4ca27aec303f54dc9ee670c67c2fefa80dd004a1.tar.gz
goflink-4ca27aec303f54dc9ee670c67c2fefa80dd004a1.zip
feat: add KeyByOperator and wordCounter for stateful processing
Diffstat (limited to 'core/operator/keyby.go')
-rw-r--r--core/operator/keyby.go32
1 files changed, 32 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
+}