From a81bd74a4ac4f1c5516d590f654afc0bb5c97b5c Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 18 Nov 2018 22:53:28 +0000 Subject: [PATCH] Use rate-limiters on DynamoDB operations --- extras/scanner/main.go | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/extras/scanner/main.go b/extras/scanner/main.go index a1f5bd15f..a534d7a42 100644 --- a/extras/scanner/main.go +++ b/extras/scanner/main.go @@ -25,6 +25,7 @@ import ( log "github.com/sirupsen/logrus" awscommon "github.com/weaveworks/common/aws" "github.com/weaveworks/common/instrument" + "golang.org/x/time/rate" ) type scanner struct { @@ -37,6 +38,9 @@ type scanner struct { bucketName string address string + writeLimiter *rate.Limiter + queryLimiter *rate.Limiter + dynamoDB *dynamodb.DynamoDB s3 *s3.S3 @@ -81,6 +85,9 @@ func main() { collectorURL string s3URL string + queryRateLimit float64 + writeRateLimit float64 + orgsFile string scanner scanner @@ -89,6 +96,8 @@ func main() { flag.StringVar(&collectorURL, "app.collector", "local", "Collector to use (local, dynamodb, or file/directory)") flag.StringVar(&s3URL, "app.collector.s3", "local", "S3 URL to use (when collector is dynamodb)") + flag.Float64Var(&queryRateLimit, "query-rate-limit", 100, "Max rate to query DynamoDB") + flag.Float64Var(&writeRateLimit, "write-rate-limit", 100, "Rate-limit on throttling from DynamoDB") flag.IntVar(&scanner.startHour, "start-hour", 406848, "Hour number to start") flag.IntVar(&scanner.stopHour, "stop-hour", 406848, "Hour number to stop (0 for current hour)") flag.IntVar(&scanner.segments, "segments", 1, "Number of segments to read in parallel") @@ -117,6 +126,9 @@ func main() { scanner.tableName = strings.TrimPrefix(parsed.Path, "/") scanner.s3 = s3.New(session.New(s3Config)) + scanner.writeLimiter = rate.NewLimiter(rate.Limit(writeRateLimit), 25) // burst size should be the largest batch + scanner.queryLimiter = rate.NewLimiter(rate.Limit(queryRateLimit), 1) // we only do one query at a time + // HTTP listener for profiling go func() { http.Handle("/metrics", promhttp.Handler()) @@ -187,10 +199,10 @@ func (sc *scanner) processOrg(ctx context.Context, org string) { for hour := sc.startHour; hour <= sc.stopHour; hour++ { var keys []map[string]*dynamodb.AttributeValue for { + sc.queryLimiter.Wait(ctx) var err error keys, err = queryDynamo(ctx, sc.dynamoDB, sc.tableName, org, int64(hour)) if throttled(err) { - time.Sleep(time.Second) continue } checkFatal(err) @@ -343,6 +355,7 @@ func (sc *scanner) deleteLoop(pending *sync.WaitGroup) { } if err != nil { if throttled(err) { + sc.writeLimiter.WaitN(context.Background(), len(batch)) // Send the whole request back into the batcher for _, item := range batch { sc.retry <- item.DeleteRequest.Key @@ -354,6 +367,9 @@ func (sc *scanner) deleteLoop(pending *sync.WaitGroup) { continue } count := 0 + if len(ret.UnprocessedItems) > 0 { + sc.writeLimiter.WaitN(context.Background(), len(ret.UnprocessedItems)) + } // Send unprocessed items back into the batcher for _, items := range ret.UnprocessedItems { count += len(items)