forked from samit22/go-ddb
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathscanner.go
97 lines (81 loc) · 2.05 KB
/
scanner.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
package ddb
import (
"expvar"
"fmt"
"sync"
"time"
"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/service/dynamodb"
"github.com/jpillora/backoff"
)
// NewScanner creates a new scanner with ddb connection
func NewScanner(config Config) *Scanner {
config.setDefaults()
return &Scanner{
waitGroup: &sync.WaitGroup{},
Config: config,
CompletedSegments: expvar.NewInt("scanner.CompletedSegments"),
}
}
// Scanner is
type Scanner struct {
waitGroup *sync.WaitGroup
Config
CompletedSegments *expvar.Int
}
// Start uses the handler function to process items for each of the total shard
func (s *Scanner) Start(handler Handler) {
for i := 0; i < s.SegmentCount; i++ {
s.waitGroup.Add(1)
segment := (s.SegmentCount * s.SegmentOffset) + i
go s.handlerLoop(handler, segment)
}
}
// Wait pauses program until waitgroup is fulfilled
func (s *Scanner) Wait() {
s.waitGroup.Wait()
}
func (s *Scanner) handlerLoop(handler Handler, segment int) {
defer s.waitGroup.Done()
var lastEvaluatedKey map[string]*dynamodb.AttributeValue
if s.Checkpoint != nil {
lastEvaluatedKey = s.Checkpoint.Get(segment)
}
bk := &backoff.Backoff{
Max: 5 * time.Minute,
Jitter: true,
}
for {
// scan params
params := &dynamodb.ScanInput{
TableName: aws.String(s.TableName),
Segment: aws.Int64(int64(segment)),
TotalSegments: aws.Int64(int64(s.TotalSegments)),
Limit: aws.Int64(s.Config.Limit),
}
// last evaluated key
if lastEvaluatedKey != nil {
params.ExclusiveStartKey = lastEvaluatedKey
}
// scan, sleep if rate limited
resp, err := s.Svc.Scan(params)
if err != nil {
fmt.Println(err)
time.Sleep(bk.Duration())
continue
}
bk.Reset()
// call the handler function with items
handler.HandleItems(resp.Items)
// exit if last evaluated key empty
if resp.LastEvaluatedKey == nil {
s.CompletedSegments.Add(1)
break
}
// set last evaluated key
lastEvaluatedKey = resp.LastEvaluatedKey
if s.Checkpoint != nil {
s.Checkpoint.Set(segment, lastEvaluatedKey)
}
}
}