Skip to content

Commit c90bc89

Browse files
authored
Support Prefix List (#2)
1 parent 421d8fd commit c90bc89

8 files changed

Lines changed: 198 additions & 19 deletions

File tree

Dockerfile

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,12 @@ ENV JOB_QUEUE_NAME ''
2222

2323
ENV SRC_BUCKET ''
2424
ENV SRC_PREFIX ''
25+
ENV SRC_PREFIX_LIST ''
2526
ENV SRC_REGION ''
2627
ENV SRC_ENDPOINT ''
2728
ENV SRC_CREDENTIALS ''
2829
ENV SRC_IN_CURRENT_ACCOUNT false
30+
ENV SKIP_COMPARE false
2931

3032
ENV DEST_BUCKET ''
3133
ENV DEST_PREFIX ''

cmd/root.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ func initConfig() {
7171
viper.SetDefault("srcType", "Amazon_S3")
7272
viper.SetDefault("destStorageClass", "STANDARD")
7373
viper.SetDefault("srcPrefix", "")
74+
viper.SetDefault("srcPrefixList", "")
7475
viper.SetDefault("srcCredential", "")
7576
viper.SetDefault("srcEndpoint", "")
7677
viper.SetDefault("destPrefix", "")
@@ -89,10 +90,12 @@ func initConfig() {
8990
viper.BindEnv("srcType", "SOURCE_TYPE")
9091
viper.BindEnv("srcBucket", "SRC_BUCKET")
9192
viper.BindEnv("srcPrefix", "SRC_PREFIX")
93+
viper.BindEnv("srcPrefixList", "SRC_PREFIX_LIST")
9294
viper.BindEnv("srcRegion", "SRC_REGION")
9395
viper.BindEnv("srcEndpoint", "SRC_ENDPOINT")
9496
viper.BindEnv("srcCredential", "SRC_CREDENTIALS")
95-
viper.BindEnv("SrcInCurrentAccount", "SRC_IN_CURRENT_ACCOUNT")
97+
viper.BindEnv("srcInCurrentAccount", "SRC_IN_CURRENT_ACCOUNT")
98+
viper.BindEnv("skipCompare", "SKIP_COMPARE")
9699

97100
viper.BindEnv("destBucket", "DEST_BUCKET")
98101
viper.BindEnv("destPrefix", "DEST_PREFIX")
@@ -146,10 +149,12 @@ func initConfig() {
146149
SrcType: viper.GetString("srcType"),
147150
SrcBucket: viper.GetString("srcBucket"),
148151
SrcPrefix: viper.GetString("srcPrefix"),
152+
SrcPrefixList: viper.GetString("srcPrefixList"),
149153
SrcRegion: viper.GetString("srcRegion"),
150154
SrcEndpoint: viper.GetString("srcEndpoint"),
151155
SrcCredential: viper.GetString("srcCredential"),
152156
SrcInCurrentAccount: viper.GetBool("srcInCurrentAccount"),
157+
SkipCompare: viper.GetBool("skipCompare"),
153158
DestBucket: viper.GetString("destBucket"),
154159
DestPrefix: viper.GetString("destPrefix"),
155160
DestRegion: viper.GetString("destRegion"),

config-example.yaml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ srcRegion: us-west-2
55
srcEndpoint:
66
srcCredential: src
77
srcInCurrentAccount: false
8+
skipCompare: false
89

910
destBucket: dest-bucket
1011
destPrefix:

dth/client.go

Lines changed: 47 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,11 @@ import (
2626
"log"
2727
"strings"
2828
"sync"
29+
"time"
2930

3031
"github.com/aws/aws-sdk-go-v2/aws"
3132
"github.com/aws/aws-sdk-go-v2/credentials"
33+
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
3234
"github.com/aws/aws-sdk-go-v2/service/s3"
3335
"github.com/aws/aws-sdk-go-v2/service/s3/types"
3436
)
@@ -42,6 +44,7 @@ type Client interface {
4244
ListCommonPrefixes(ctx context.Context, depth int, maxKeys int32) (prefixes []*string)
4345
ListParts(ctx context.Context, key, uploadID *string) (parts map[int]*Part)
4446
GetUploadID(ctx context.Context, key *string) (uploadID *string)
47+
ListSelectedPrefixes(ctx context.Context, key *string) (prefixes []*string)
4548

4649
// WRITE
4750
PutObject(ctx context.Context, key *string, body []byte, storageClass, acl *string, meta *Metadata) (etag *string, err error)
@@ -54,8 +57,8 @@ type Client interface {
5457

5558
// S3Client is an implementation of Client interface for Amazon S3
5659
type S3Client struct {
57-
bucket, prefix, region, sourceType string
58-
client *s3.Client
60+
bucket, prefix, prefixList, region, sourceType string
61+
client *s3.Client
5962
}
6063

6164
// S3Credentials is
@@ -82,7 +85,7 @@ func getEndpointURL(region, sourceType string) (url string) {
8285
}
8386

8487
// NewS3Client creates a S3Client instance
85-
func NewS3Client(ctx context.Context, bucket, prefix, endpoint, region, sourceType string, cred *S3Credentials) *S3Client {
88+
func NewS3Client(ctx context.Context, bucket, prefix, prefixList, endpoint, region, sourceType string, cred *S3Credentials) *S3Client {
8689

8790
cfg := loadDefaultConfig(ctx)
8891

@@ -122,6 +125,7 @@ func NewS3Client(ctx context.Context, bucket, prefix, endpoint, region, sourceTy
122125
return &S3Client{
123126
bucket: bucket,
124127
prefix: prefix,
128+
prefixList: prefixList,
125129
client: client,
126130
region: region,
127131
sourceType: sourceType,
@@ -311,6 +315,46 @@ func (c *S3Client) HeadObject(ctx context.Context, key *string) *Metadata {
311315

312316
}
313317

318+
// ListSelectedPrefixes is a function to list prefixes from a customized list file.
319+
func (c *S3Client) ListSelectedPrefixes(ctx context.Context, key *string) (prefixes []*string) {
320+
321+
downloader := manager.NewDownloader(c.client)
322+
323+
getBuf := manager.NewWriteAtBuffer([]byte{})
324+
325+
input := &s3.GetObjectInput{
326+
Bucket: &c.bucket,
327+
Key: key,
328+
}
329+
330+
dounload_start := time.Now()
331+
log.Printf("Start downloading the Prefix List File.")
332+
_, err := downloader.Download(ctx, getBuf, input)
333+
download_end := time.Since(dounload_start)
334+
if err != nil {
335+
fmt.Print(err)
336+
} else {
337+
log.Printf("Download the Prefix List File Completed in %v\n", download_end)
338+
}
339+
start := time.Now()
340+
prefixes_value := make([]string, 0, 100000000)
341+
342+
for i, m := range strings.Split(string(getBuf.Bytes()), "\n") {
343+
if i > 100000000 {
344+
log.Printf("The number of prefixes in the list file is larger than 100,000,000, please seperate the file.")
345+
return
346+
}
347+
if len(m) > 0 {
348+
prefixes_value = append(prefixes_value, m)
349+
prefixes = append(prefixes, &prefixes_value[i])
350+
}
351+
}
352+
log.Printf("Got %d prefixes from the customized list file.", len(prefixes))
353+
end := time.Since(start)
354+
log.Printf("Getting Prefixes List Job Completed in %v\n", end)
355+
return
356+
}
357+
314358
// PutObject is a function to put (upload) an object to Amazon S3
315359
func (c *S3Client) PutObject(ctx context.Context, key *string, body []byte, storageClass, acl *string, meta *Metadata) (etag *string, err error) {
316360
// log.Printf("S3> Uploading object %s to bucket %s\n", key, c.bucket)

dth/config.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -66,10 +66,10 @@ type JobOptions struct {
6666

6767
// JobConfig is General Job Info
6868
type JobConfig struct {
69-
SrcType, SrcBucket, SrcPrefix, SrcRegion, SrcEndpoint, SrcCredential string
70-
DestBucket, DestPrefix, DestRegion, DestCredential, DestStorageClass, DestAcl string
71-
JobTableName, JobQueueName string
72-
SrcInCurrentAccount, DestInCurrentAccount bool
69+
SrcType, SrcBucket, SrcPrefix, SrcPrefixList, SrcRegion, SrcEndpoint, SrcCredential string
70+
DestBucket, DestPrefix, DestRegion, DestCredential, DestStorageClass, DestAcl string
71+
JobTableName, JobQueueName string
72+
SrcInCurrentAccount, DestInCurrentAccount, SkipCompare bool
7373
*JobOptions
7474
}
7575

dth/job.go

Lines changed: 105 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -113,8 +113,8 @@ func NewFinder(ctx context.Context, cfg *JobConfig) (f *Finder) {
113113
srcCred := getCredentials(ctx, cfg.SrcCredential, cfg.SrcInCurrentAccount, sm)
114114
desCred := getCredentials(ctx, cfg.DestCredential, cfg.DestInCurrentAccount, sm)
115115

116-
srcClient := NewS3Client(ctx, cfg.SrcBucket, cfg.SrcPrefix, cfg.SrcEndpoint, cfg.SrcRegion, cfg.SrcType, srcCred)
117-
desClient := NewS3Client(ctx, cfg.DestBucket, cfg.DestPrefix, "", cfg.DestRegion, "Amazon_S3", desCred)
116+
srcClient := NewS3Client(ctx, cfg.SrcBucket, cfg.SrcPrefix, cfg.SrcPrefixList, cfg.SrcEndpoint, cfg.SrcRegion, cfg.SrcType, srcCred)
117+
desClient := NewS3Client(ctx, cfg.DestBucket, cfg.DestPrefix, "", "", cfg.DestRegion, "Amazon_S3", desCred)
118118

119119
f = &Finder{
120120
srcClient: srcClient,
@@ -147,16 +147,27 @@ func (f *Finder) Run(ctx context.Context) {
147147
// Note that bigger number needs more memory
148148
compareCh := make(chan struct{}, f.cfg.FinderNumber)
149149

150-
prefixes := f.srcClient.ListCommonPrefixes(ctx, f.cfg.FinderDepth, f.cfg.MaxKeys)
150+
var prefixes []*string
151+
log.Printf("Prefix List File: %s", f.cfg.SrcPrefixList)
151152

153+
if len(f.cfg.SrcPrefixList) > 0 {
154+
prefixes = f.srcClient.ListSelectedPrefixes(ctx, &f.cfg.SrcPrefixList)
155+
} else {
156+
prefixes = f.srcClient.ListCommonPrefixes(ctx, f.cfg.FinderDepth, f.cfg.MaxKeys)
157+
}
152158
var wg sync.WaitGroup
153159

154160
start := time.Now()
155161

156162
for _, p := range prefixes {
157163
compareCh <- struct{}{}
164+
log.Printf("prefix: %s", *p)
158165
wg.Add(1)
159-
go f.compareAndSend(ctx, p, batchCh, msgCh, compareCh, &wg)
166+
if f.cfg.SkipCompare {
167+
go f.directSend(ctx, p, batchCh, msgCh, compareCh, &wg)
168+
} else {
169+
go f.compareAndSend(ctx, p, batchCh, msgCh, compareCh, &wg)
170+
}
160171
}
161172
wg.Wait()
162173

@@ -295,6 +306,94 @@ func (f *Finder) compareAndSend(ctx context.Context, prefix *string, batchCh cha
295306
<-compareCh
296307
}
297308

309+
// This function will send the task to SQS Queue directly, without comparison.
310+
func (f *Finder) directSend(ctx context.Context, prefix *string, batchCh chan struct{}, msgCh chan *string, compareCh chan struct{}, wg *sync.WaitGroup) {
311+
defer wg.Done()
312+
313+
log.Printf("Scanning prefix /%s\n", *prefix)
314+
315+
token := ""
316+
i, j := 0, 0
317+
retry := 0
318+
// batch := make([]*string, f.cfg.MessageBatchSize)
319+
320+
log.Printf("Start sending without comparison ...\n")
321+
// start := time.Now()
322+
323+
for token != "End" {
324+
// source := f.getSourceObjects(ctx, &token, prefix)
325+
source, err := f.srcClient.ListObjects(ctx, &token, prefix, f.cfg.MaxKeys)
326+
if err != nil {
327+
log.Printf("Fail to get source list - %s\n", err.Error())
328+
//
329+
log.Printf("Sleep for 1 minute and try again...")
330+
retry++
331+
332+
if retry <= MaxRetries {
333+
time.Sleep(time.Minute * 1)
334+
continue
335+
} else {
336+
log.Printf("Still unable to list source list after %d retries\n", MaxRetries)
337+
// Log the last token and exit
338+
log.Fatalf("The last token is %s\n", token)
339+
}
340+
341+
}
342+
343+
// if a successful list, reset to 0
344+
retry = 0
345+
346+
for _, obj := range source {
347+
// TODO: Check if there is another way to compare
348+
// Currently, map is used to search if such object exists in target
349+
msgCh <- obj.toString()
350+
i++
351+
if i%f.cfg.MessageBatchSize == 0 {
352+
wg.Add(1)
353+
j++
354+
if j%100 == 0 {
355+
log.Printf("Found %d batches in prefix /%s\n", j, *prefix)
356+
}
357+
batchCh <- struct{}{}
358+
359+
// start a go routine to send messages in batch
360+
go func(i int) {
361+
defer wg.Done()
362+
batch := make([]*string, i)
363+
for a := 0; a < i; a++ {
364+
batch[a] = <-msgCh
365+
}
366+
367+
f.sqs.SendMessageInBatch(ctx, batch)
368+
<-batchCh
369+
}(f.cfg.MessageBatchSize)
370+
i = 0
371+
}
372+
}
373+
}
374+
// For remainning objects.
375+
if i != 0 {
376+
j++
377+
wg.Add(1)
378+
batchCh <- struct{}{}
379+
go func(i int) {
380+
defer wg.Done()
381+
batch := make([]*string, i)
382+
for a := 0; a < i; a++ {
383+
batch[a] = <-msgCh
384+
}
385+
386+
f.sqs.SendMessageInBatch(ctx, batch)
387+
<-batchCh
388+
}(i)
389+
}
390+
391+
// end := time.Since(start)
392+
// log.Printf("Compared and Sent %d batches in %v", j, end)
393+
log.Printf("Completed in prefix /%s, found %d batches in total", *prefix, j)
394+
<-compareCh
395+
}
396+
298397
// NewWorker creates a new Worker instance
299398
func NewWorker(ctx context.Context, cfg *JobConfig) (w *Worker) {
300399
log.Printf("Source Type is %s\n", cfg.SrcType)
@@ -310,8 +409,8 @@ func NewWorker(ctx context.Context, cfg *JobConfig) (w *Worker) {
310409
srcCred := getCredentials(ctx, cfg.SrcCredential, cfg.SrcInCurrentAccount, sm)
311410
desCred := getCredentials(ctx, cfg.DestCredential, cfg.DestInCurrentAccount, sm)
312411

313-
srcClient := NewS3Client(ctx, cfg.SrcBucket, cfg.SrcPrefix, cfg.SrcEndpoint, cfg.SrcRegion, cfg.SrcType, srcCred)
314-
desClient := NewS3Client(ctx, cfg.DestBucket, cfg.DestPrefix, "", cfg.DestRegion, "Amazon_S3", desCred)
412+
srcClient := NewS3Client(ctx, cfg.SrcBucket, cfg.SrcPrefix, cfg.SrcPrefixList, cfg.SrcEndpoint, cfg.SrcRegion, cfg.SrcType, srcCred)
413+
desClient := NewS3Client(ctx, cfg.DestBucket, cfg.DestPrefix, "", "", cfg.DestRegion, "Amazon_S3", desCred)
315414

316415
return &Worker{
317416
srcClient: srcClient,

go.mod

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,15 +3,16 @@ module golang.a2z.com/dthcli
33
go 1.16
44

55
require (
6-
github.com/aws/aws-sdk-go-v2 v1.6.0
7-
github.com/aws/aws-sdk-go-v2/config v1.3.0
8-
github.com/aws/aws-sdk-go-v2/credentials v1.2.1
6+
github.com/aws/aws-sdk-go-v2 v1.8.1
7+
github.com/aws/aws-sdk-go-v2/config v1.6.1
8+
github.com/aws/aws-sdk-go-v2/credentials v1.3.3
99
github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue v1.1.1
10+
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.4.1
1011
github.com/aws/aws-sdk-go-v2/service/dynamodb v1.3.1
11-
github.com/aws/aws-sdk-go-v2/service/s3 v1.9.0
12+
github.com/aws/aws-sdk-go-v2/service/s3 v1.13.0
1213
github.com/aws/aws-sdk-go-v2/service/secretsmanager v1.3.1
1314
github.com/aws/aws-sdk-go-v2/service/sqs v1.4.1
14-
github.com/aws/smithy-go v1.4.0
15+
github.com/aws/smithy-go v1.7.0
1516
github.com/spf13/cobra v1.1.3
1617
github.com/spf13/viper v1.7.1
1718
)

0 commit comments

Comments
 (0)