Repository navigation
Expand file tree
/
Copy pathparallel.go
More file actions
executable file
·85 lines (67 loc) · 1.58 KB
/
Copy pathparallel.go
File metadata and controls
executable file
·85 lines (67 loc) · 1.58 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
package co
import (
"sync"
"sync/atomic"
"go.tempura.ink/co/ds/pool"
syncx "go.tempura.ink/co/internal/syncx"
)
type parallel[R any] struct {
workerPool *pool.WorkerPool[R]
seqFnMap sync.Map
ifPersistentData bool
storedData *List[R]
sent uint64
finished uint64
recieverCond *syncx.Condx
}
func NewParallel[R any](maxWorkers int) *parallel[R] {
d := ¶llel[R]{
workerPool: pool.NewWorkerPool[R](maxWorkers),
storedData: NewList[R](),
recieverCond: syncx.NewCondx(&sync.Mutex{}),
}
d.workerPool.SetCallbackFn(d.receiveValue)
return d
}
func (d *parallel[R]) SetPersistentData(b bool) *parallel[R] {
d.ifPersistentData = b
return d
}
func (d *parallel[R]) Process(fn func() R) chan R {
atomic.AddUint64(&d.sent, 1)
seq := d.workerPool.ReserveSeq()
ch := make(chan R)
d.seqFnMap.Store(seq, ch)
d.workerPool.AddJobAt(seq, fn)
return ch
}
func (d *parallel[R]) receiveValue(seq uint64, val R) {
ch, ok := d.seqFnMap.Load(seq)
if ok {
syncx.SafeNSend(ch.(chan R), val)
}
d.seqFnMap.Delete(seq)
if d.ifPersistentData {
d.storedData.setAt(int(seq)-1, val)
}
d.recieverCond.Broadcastify(&syncx.BroadcastOption{
PreProcessFn: func() {
atomic.AddUint64(&d.finished, 1)
}},
)
}
func (d *parallel[R]) GetData() []R {
if !d.ifPersistentData {
panic("co/parallel error when get data: persistent data mode is not set")
}
return d.storedData.list
}
func (d *parallel[R]) Wait() *parallel[R] {
d.workerPool.Wait()
d.recieverCond.Waitify(&syncx.WaitOption{
ConditionFn: func() bool {
return d.finished != d.sent
},
})
return d
}