diff options
| author | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-09 18:32:44 +0700 |
|---|---|---|
| committer | SangTran-127 <tranquangsang12.7@gmail.com> | 2026-08-09 18:32:44 +0700 |
| commit | 4ca27aec303f54dc9ee670c67c2fefa80dd004a1 (patch) | |
| tree | 5628714f47e374bfb98c1cea99970b19892e8b3e /core/operator/keyby.go | |
| parent | a892cb5b1d121177881fa56a79abbafef5880a80 (diff) | |
| download | goflink-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.go | 32 |
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 +} |