Repository navigation
Expand file tree
/
Copy pathkey_join_group_async.go
More file actions
90 lines (78 loc) · 1.8 KB
/
Copy pathkey_join_group_async.go
File metadata and controls
90 lines (78 loc) · 1.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
package flame
import (
"sync"
"golang.org/x/exp/constraints"
)
type KeyJoinGroupAsyncNode[K constraints.Ordered, X, Z any] struct {
Inputs []chan KeyValue[K, X]
Outputs []chan KeyValue[K, Z]
Proc func(K, []X) Z
}
func AddKeyJoinGroupAsync[K constraints.Ordered, X, Z any](w *Workflow, f func(K, []X) Z) *KeyJoinGroupAsyncNode[K, X, Z] {
n := &KeyJoinGroupAsyncNode[K, X, Z]{Proc: f, Outputs: []chan KeyValue[K, Z]{}}
w.Nodes = append(w.Nodes, n)
return n
}
func (n *KeyJoinGroupAsyncNode[K, X, Z]) start(wf *Workflow) {
wf.WaitGroup.Add(1)
mut := &sync.Mutex{}
store := map[K][]X{}
updates := make(chan K, len(n.Inputs))
wg := &sync.WaitGroup{}
for i := 0; i < len(n.Inputs); i++ {
wg.Add(1)
go func(inputN int) {
for l := range n.Inputs[inputN] {
mut.Lock()
if v, ok := store[l.Key]; ok {
v[inputN] = l.Value
} else {
v := make([]X, len(n.Inputs))
v[inputN] = l.Value
store[l.Key] = v
}
mut.Unlock()
updates <- l.Key
}
wg.Done()
}(i)
}
go func() {
wg.Wait()
close(updates)
}()
go func() {
hits := map[K]int{}
for key := range updates {
if x, ok := hits[key]; ok {
x++
hits[key] = x
} else {
hits[key] = 1
}
if hits[key] >= len(n.Inputs) {
delete(hits, key)
mut.Lock()
val := n.Proc(key, store[key])
delete(store, key)
mut.Unlock()
for i := range n.Outputs {
n.Outputs[i] <- KeyValue[K, Z]{key, val}
}
}
}
for i := range n.Outputs {
close(n.Outputs[i])
}
wf.WaitGroup.Done()
}()
}
func (n *KeyJoinGroupAsyncNode[K, X, Z]) GetOutput() chan KeyValue[K, Z] {
m := make(chan KeyValue[K, Z])
n.Outputs = append(n.Outputs, m)
return m
}
func (n *KeyJoinGroupAsyncNode[K, X, Z]) Connect(e Emitter[KeyValue[K, X]]) {
o := e.GetOutput()
n.Inputs = append(n.Inputs, o)
}