-
Notifications
You must be signed in to change notification settings - Fork 3
/
key_scanner.go
137 lines (123 loc) · 2.85 KB
/
key_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
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
package asc
import (
"bytes"
"fmt"
"github.com/aerospike/aerospike-client-go"
"github.com/viant/toolbox"
"io"
"log"
"os"
"path"
"sync"
"sync/atomic"
)
const bufferSize = 1024 * 1024 * 256
type KeyScanner struct {
client *aerospike.Client
scanPolicy *aerospike.ScanPolicy
baseDirectory string
namespace string
dataSet string
writer *bufferedWriter
}
func (s *KeyScanner) scanAll() ([]string, error) {
filename := path.Join(s.baseDirectory, fmt.Sprintf("%v-%v.keys", s.namespace, s.dataSet))
toolbox.RemoveFileIfExist(filename)
writer, err := os.OpenFile(filename, os.O_CREATE|os.O_WRONLY, 0644)
if err != nil {
return nil, err
}
s.writer = newBufferedWriter(writer, bufferSize)
defer s.writer.Close()
var result = []string{filename}
recordSet, err := s.client.ScanAll(s.scanPolicy, s.namespace, s.dataSet)
if err != nil {
return nil, err
}
records := recordSet.Records
errors := recordSet.Errors
for recordSet.IsActive() {
select {
case record := <-records:
if record != nil && record.Key != nil {
err = WriteKey(record.Key, s.writer)
if err != nil {
return nil, err
}
}
case err = <-errors:
if err != nil {
log.Printf("failed to scan keys: %v", err)
return nil, err
}
}
}
return result, nil
}
func (s *KeyScanner) Scan() ([]string, error) {
return s.scanAll()
}
func NewKeyScanner(client *aerospike.Client, scanPolicy *aerospike.ScanPolicy, baseDirectory, namespace, dataSet string) *KeyScanner {
return &KeyScanner{
client: client,
scanPolicy: scanPolicy,
baseDirectory: baseDirectory,
namespace: namespace,
dataSet: dataSet,
}
}
//writer tries to hold in memory keys
type bufferedWriter struct {
size int32
bufSize int
flushCount int32
buf *bytes.Buffer
writer io.WriteCloser
mux *sync.Mutex
}
func (w *bufferedWriter) flush() error {
var data = w.buf.Bytes()
written, err := w.writer.Write(data)
if err != nil {
return err
}
if written != len(data) {
return fmt.Errorf("wrote %v of %v", written, len(data))
}
w.buf.Reset()
atomic.AddInt32(&w.flushCount, 1)
atomic.StoreInt32(&w.size, 0)
return nil
}
func (w *bufferedWriter) pendingSize() int {
return int(atomic.LoadInt32(&w.size))
}
func (w *bufferedWriter) Write(p []byte) (n int, err error) {
w.mux.Lock()
defer w.mux.Unlock()
if w.pendingSize()+len(p) > w.bufSize {
err = w.flush()
if err != nil {
return 0, err
}
}
atomic.AddInt32(&w.size, int32(len(p)))
return w.buf.Write(p)
}
func (w *bufferedWriter) Close() (err error) {
if w.pendingSize() > 0 {
err = w.flush()
if err != nil {
return err
}
}
return w.writer.Close()
}
func newBufferedWriter(writer io.WriteCloser, bufSize int) *bufferedWriter {
return &bufferedWriter{
bufSize: bufSize,
buf: new(bytes.Buffer),
writer: writer,
mux: &sync.Mutex{},
}
}