-
Notifications
You must be signed in to change notification settings - Fork 0
/
db_find_batch.go
119 lines (104 loc) · 2.38 KB
/
db_find_batch.go
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
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
package do
import (
"context"
"fmt"
)
// Batch process data find from Storer in batches
func Batch[S Storer, F Finder[R], R any](db S, finder F, batchNum int, handler func([]R) error) (err error) {
query, args := finder.Query()
rows, err := db.QueryContext(context.TODO(), query, args...)
if err != nil {
return
}
defer rows.Close()
colTypes, err := rows.ColumnTypes()
if err != nil {
return
}
batch := make([]R, 0, batchNum)
for rows.Next() {
t, fields := finder.NewScanObjAndFields(colTypes)
if err = rows.Scan(fields...); err != nil {
return
}
batch = append(batch, *t)
if batchNum > 0 && len(batch) >= batchNum {
if err = handler(batch); err != nil {
err = fmt.Errorf("batch handle failed %w", err)
return
}
batch = make([]R, 0, batchNum)
}
}
if err = rows.Err(); err != nil {
return
}
if len(batch) != 0 {
if err = handler(batch); err != nil {
err = fmt.Errorf("batch handle failed %w", err)
return
}
}
return
}
// BatchConcurrent batch process data concurrently
func BatchConcurrent[S Storer, F Finder[R], R any](db S, finder F, batchNum int, handler func([]R) error, concNum int) (err error) {
query, args := finder.Query()
rows, err := db.QueryContext(context.TODO(), query, args...)
if err != nil {
return
}
defer rows.Close()
colTypes, err := rows.ColumnTypes()
if err != nil {
return
}
batchWorker := NewWorker(concNum)
batchWorker.Start()
defer batchWorker.Stop()
batch := make([]R, 0, batchNum)
for rows.Next() {
t, fields := finder.NewScanObjAndFields(colTypes)
if err = rows.Scan(fields...); err != nil {
return
}
batch = append(batch, *t)
if batchNum > 0 && len(batch) >= batchNum {
batchCopy := batch
if err = batchWorker.Push(*NewJob(
DoWithCtx(func(ctx context.Context) error {
if err = handler(batchCopy); err != nil {
err = fmt.Errorf("batch handle failed %w", err)
return err
}
return nil
}),
0,
nil,
)); err != nil {
return
}
batch = make([]R, 0, batchNum)
}
}
if err = rows.Err(); err != nil {
return
}
if len(batch) != 0 {
batchCopy := batch
if err = batchWorker.Push(*NewJob(
DoWithCtx(func(ctx context.Context) error {
if err = handler(batchCopy); err != nil {
err = fmt.Errorf("batch handle failed %w", err)
return err
}
return nil
}),
0,
nil,
)); err != nil {
return
}
}
return
}