This repository was archived by the owner on Jun 24, 2026. It is now read-only.
forked from katanomi/pkg
-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathpage.go
More file actions
91 lines (73 loc) · 2.35 KB
/
Copy pathpage.go
File metadata and controls
91 lines (73 loc) · 2.35 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
91
/*
Copyright 2022 The AlaudaDevops Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package parallel
import (
"context"
"fmt"
"time"
"k8s.io/utils/trace"
"knative.dev/pkg/logging"
)
// PageRequestFunc is a tool for concurrent processing of pagination
type PageRequestFunc struct {
// RequestPage for concurrent request paging
RequestPage func(ctx context.Context, pageSize int, page int) (interface{}, error)
// PageResult for get paging information
PageResult func(items interface{}) (total int, currentPageLen int, err error)
}
// PageRequest is concurrent request paging
func PageRequest(ctx context.Context, logName string, concurrency int, pageSize int, f PageRequestFunc) ([]interface{}, error) {
log := trace.New("PageRequest", trace.Field{Key: "name", Value: logName})
logger := logging.FromContext(ctx)
defer func() {
log.LogIfLong(3 * time.Second)
}()
items, err := f.RequestPage(ctx, pageSize, 1)
if err != nil {
return nil, err
}
log.Step("requested page 1")
total, firstPageLen, err := f.PageResult(items)
if err != nil {
return nil, err
}
if firstPageLen < pageSize {
return []interface{}{items}, nil
}
if total == firstPageLen {
return []interface{}{items}, nil
}
var request = func(i int) func() (interface{}, error) {
return func() (interface{}, error) {
items, err := f.RequestPage(ctx, pageSize, i)
log.Step(fmt.Sprintf("requested page %d", i))
return items, err
}
}
totalPage := total / pageSize
if total%pageSize != 0 {
totalPage = totalPage + 1
}
if totalPage-1 < concurrency { // first page we have requested, so skip first page
concurrency = totalPage - 1
}
p := P(logger, "PageRequest").FailFast().SetConcurrent(concurrency).Context(ctx)
for i := 2; i <= totalPage; i++ {
p.Add(request(i))
}
results, err := p.Do().Wait()
if err != nil {
return nil, err
}
return append([]interface{}{items}, results...), nil
}