From 777ff07e1938bd54354a0712ee40f0a8e094abf0 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Thu, 2 Apr 2020 21:13:25 +0000 Subject: [PATCH 1/6] refactor(multitenant): break report storage code out into sub-functions So the main Add() function isn't so long. --- app/multitenant/aws_collector.go | 116 ++++++++++++++++--------------- 1 file changed, 59 insertions(+), 57 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 0e84bdba1..25c80dcdb 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -160,6 +160,9 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { registerAWSCollectorMetricsOnce.Do(registerAWSCollectorMetrics) var nc *nats.Conn if config.NatsHost != "" { + if config.MemcacheClient == nil { + return nil, fmt.Errorf("Must supply memcache client when using nats") + } var err error nc, err = nats.Connect(config.NatsHost) if err != nil { @@ -481,13 +484,46 @@ func calculateDynamoKeys(userid string, now time.Time) (string, string) { return rowKey, colKey } -// calculateReportKey determines the key we should use for a report. -func calculateReportKey(rowKey, colKey string) (string, error) { +// calculateReportKeys returns DynamoDB row & col keys, and S3/memcached key that we will use for a report +func calculateReportKeys(userid string, now time.Time) (string, string, string) { + rowKey, colKey := calculateDynamoKeys(userid, now) rowKeyHash := md5.New() - if _, err := io.WriteString(rowKeyHash, rowKey); err != nil { - return "", err + _, _ = io.WriteString(rowKeyHash, rowKey) // hash write doesn't error + return rowKey, colKey, fmt.Sprintf("%x/%s", rowKeyHash.Sum(nil), colKey) +} + +func (c *awsCollector) persistReport(ctx context.Context, userid, rowKey, colKey, reportKey string, buf []byte) error { + // Put in S3 and cache before index, so it is fetchable before it is discoverable + reportSize, err := c.cfg.S3Store.StoreReportBytes(ctx, reportKey, buf) + if err != nil { + return err } - return fmt.Sprintf("%x/%s", rowKeyHash.Sum(nil), colKey), nil + if c.cfg.MemcacheClient != nil { + _, err = c.cfg.MemcacheClient.StoreReportBytes(ctx, reportKey, buf) + if err != nil { + log.Warningf("Could not store %v in memcache: %v", reportKey, err) + } + } + + dynamoValueSize.WithLabelValues("PutItem").Add(float64(len(reportKey))) + + err = instrument.TimeRequestHistogram(ctx, "DynamoDB.PutItem", dynamoRequestDuration, func(_ context.Context) error { + resp, err := c.putItemInDynamo(rowKey, colKey, reportKey) + if resp.ConsumedCapacity != nil { + dynamoConsumedCapacity.WithLabelValues("PutItem"). + Add(float64(*resp.ConsumedCapacity.CapacityUnits)) + } + return err + }) + if err != nil { + return err + } + + reportSizeHistogram.Observe(float64(reportSize)) + reportSizePerUser.WithLabelValues(userid).Add(float64(reportSize)) + reportsPerUser.WithLabelValues(userid).Inc() + + return nil } func (c *awsCollector) putItemInDynamo(rowKey, colKey, reportKey string) (*dynamodb.PutItemOutput, error) { @@ -537,63 +573,29 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e return err } - // first, put the report on s3 - rowKey, colKey := calculateDynamoKeys(userid, time.Now()) - reportKey, err := calculateReportKey(rowKey, colKey) - if err != nil { - return err - } - // Shortcut reports are published to nats but not persisted - // we'll get a full report from the same probe in a few seconds - if !rep.Shortcut { - reportSize, err := c.cfg.S3Store.StoreReportBytes(ctx, reportKey, buf) - if err != nil { - return err + if rep.Shortcut { + if c.nats != nil { + _, _, reportKey := calculateReportKeys(userid, time.Now()) + _, err = c.cfg.MemcacheClient.StoreReportBytes(ctx, reportKey, buf) + if err != nil { + log.Warningf("Could not store %v in memcache: %v", reportKey, err) + return nil + } + err := c.nats.Publish(userid, []byte(reportKey)) + natsRequests.WithLabelValues("Publish", instrument.ErrorCode(err)).Add(1) + if err != nil { + log.Errorf("Error sending shortcut report: %v", err) + } } - reportSizeHistogram.Observe(float64(reportSize)) - reportSizePerUser.WithLabelValues(userid).Add(float64(reportSize)) - reportsPerUser.WithLabelValues(userid).Inc() + return nil } - // third, put it in memcache - if c.cfg.MemcacheClient != nil { - _, err = c.cfg.MemcacheClient.StoreReportBytes(ctx, reportKey, buf) - if err != nil { - // NOTE: We don't abort here because failing to store in memcache - // doesn't actually break anything else -- it's just an - // optimization. - log.Warningf("Could not store %v in memcache: %v", reportKey, err) - } - } - - if !rep.Shortcut { - // fourth, put the key in dynamodb - dynamoValueSize.WithLabelValues("PutItem"). - Add(float64(len(reportKey))) - - var resp *dynamodb.PutItemOutput - err = instrument.TimeRequestHistogram(ctx, "DynamoDB.PutItem", dynamoRequestDuration, func(_ context.Context) error { - var err error - resp, err = c.putItemInDynamo(rowKey, colKey, reportKey) - return err - }) - - if resp.ConsumedCapacity != nil { - dynamoConsumedCapacity.WithLabelValues("PutItem"). - Add(float64(*resp.ConsumedCapacity.CapacityUnits)) - } - if err != nil { - return err - } - } - - if rep.Shortcut && c.nats != nil { - err := c.nats.Publish(userid, []byte(reportKey)) - natsRequests.WithLabelValues("Publish", instrument.ErrorCode(err)).Add(1) - if err != nil { - log.Errorf("Error sending shortcut report: %v", err) - } + rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) + err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf) + if err != nil { + return err } return nil From 104b9cba50340e037bc851c4cac34d17d88bf787 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 13 Apr 2020 09:29:19 +0000 Subject: [PATCH 2/6] refactor: Call Close() on collector Doesn't do anything at present, but will be used later. Change the signature on BillingEmitter.Close() to match. Note we didn't use the error returned. --- app/collector.go | 7 +++++++ app/multitenant/aws_collector.go | 4 ++++ app/multitenant/billing_emitter.go | 5 +++-- prog/app.go | 2 +- 4 files changed, 15 insertions(+), 3 deletions(-) diff --git a/app/collector.go b/app/collector.go index f59f6f75f..d7504da62 100644 --- a/app/collector.go +++ b/app/collector.go @@ -58,6 +58,7 @@ type Adder interface { type Collector interface { Reporter Adder + Close() } // Collector receives published reports from multiple producers. It yields a @@ -112,6 +113,9 @@ func NewCollector(window time.Duration) Collector { } } +// Close is a no-op for the regular collector +func (c *collector) Close() {} + // Add adds a report to the collector's internal state. It implements Adder. func (c *collector) Add(_ context.Context, rpt report.Report, _ []byte) error { c.mtx.Lock() @@ -243,6 +247,9 @@ func (c StaticCollector) Report(context.Context, time.Time) (report.Report, erro return report.Report(c), nil } +// Close is a no-op for the static collector +func (c StaticCollector) Close() {} + // HasReports indicates whether the collector contains reports between // timestamp-app.window and timestamp. func (c StaticCollector) HasReports(context.Context, time.Time) (bool, error) { diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 25c80dcdb..1d3396806 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -182,6 +182,10 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { }, nil } +// Close is a no-op for awsCollector +func (c *awsCollector) Close() { +} + // CreateTables creates the required tables in dynamodb func (c *awsCollector) CreateTables() error { // see if tableName exists diff --git a/app/multitenant/billing_emitter.go b/app/multitenant/billing_emitter.go index aca58b784..d58fe9e68 100644 --- a/app/multitenant/billing_emitter.go +++ b/app/multitenant/billing_emitter.go @@ -166,6 +166,7 @@ func hasWeaveNet(r report.Report) bool { } // Close shuts down the billing emitter and billing client flushing events. -func (e *BillingEmitter) Close() error { - return e.billing.Close() +func (e *BillingEmitter) Close() { + e.Collector.Close() + _ = e.billing.Close() } diff --git a/prog/app.go b/prog/app.go index e9b512b87..7e08fdb83 100644 --- a/prog/app.go +++ b/prog/app.go @@ -259,9 +259,9 @@ func appMain(flags appFlags) { log.Fatalf("Error creating emitter: %v", err) return } - defer billingEmitter.Close() collector = billingEmitter } + defer collector.Close() controlRouter, err := controlRouterFactory(userIDer, flags.controlRouterURL, flags.controlRPCTimeout) if err != nil { From ccf031b8a96bf0e06f65bff4a9a3ac3c23aa7534 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 13 Apr 2020 10:43:07 +0000 Subject: [PATCH 3/6] enhancement(multitenant): merge incoming reports in a time window This means we store fewer, bigger, reports, which reduces cost of storage and time to render when data is viewed. --- app/multitenant/aws_collector.go | 77 +++++++++++++++++++++++++++++--- prog/app.go | 5 ++- prog/main.go | 2 + 3 files changed, 75 insertions(+), 9 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 1d3396806..9179bca86 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -123,17 +123,27 @@ type AWSCollectorConfig struct { DynamoDBConfig *aws.Config DynamoTable string S3Store *S3Store + StoreInterval time.Duration NatsHost string MemcacheClient *MemcacheClient Window time.Duration MaxTopNodes int } +// if StoreInterval is set, reports are merged into here and held until flushed to store +type pendingEntry struct { + sync.Mutex + report report.Report + count int +} + type awsCollector struct { cfg AWSCollectorConfig db *dynamodb.DynamoDB merger app.Merger inProcess inProcessStore + pending sync.Map + ticker *time.Ticker nats *nats.Conn waitersLock sync.Mutex @@ -172,18 +182,60 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { // (window * report rate) * number of hosts per user * number of users reportCacheSize := (int(config.Window.Seconds()) / 3) * 10 * 5 - return &awsCollector{ + c := &awsCollector{ cfg: config, db: dynamodb.New(session.New(config.DynamoDBConfig)), merger: app.NewFastMerger(), inProcess: newInProcessStore(reportCacheSize, config.Window+reportQuantisationInterval), nats: nc, waiters: map[watchKey]*nats.Subscription{}, - }, nil + } + + if config.StoreInterval != 0 { + c.ticker = time.NewTicker(config.StoreInterval) + go c.flushLoop() + } + return c, nil } -// Close is a no-op for awsCollector +func (c *awsCollector) flushLoop() { + for _ = range c.ticker.C { + c.flushPending(context.Background()) + } +} + +// Range over all users (instances) that have pending reports and send to store +func (c *awsCollector) flushPending(ctx context.Context) { + c.pending.Range(func(key, value interface{}) bool { + userid := key.(string) + entry := value.(*pendingEntry) + + entry.Lock() + rpt, count := entry.report, entry.count + entry.report, entry.count = report.MakeReport(), 0 + entry.Unlock() + + if count > 0 { + buf, err := rpt.WriteBinary() + if err != nil { + log.Errorf("Could not serialise combined report: %v", err) + return true + } + rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) + err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf.Bytes()) + if err != nil { + log.Errorf("Could not persist combined report: %v", err) + return true + } + } + return true + }) +} + +// Close will flush pending data func (c *awsCollector) Close() { + c.ticker.Stop() // note this doesn't close the chan; goroutine keeps running + c.flushPending(context.Background()) } // CreateTables creates the required tables in dynamodb @@ -596,10 +648,21 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e return nil } - rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) - err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf) - if err != nil { - return err + if c.cfg.StoreInterval == 0 { + rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) + err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf) + if err != nil { + return err + } + } else { + entry := &pendingEntry{report: report.MakeReport()} + if e, found := c.pending.LoadOrStore(userid, entry); found { + entry = e.(*pendingEntry) + } + entry.Lock() + entry.report.UnsafeMerge(rep) + entry.count++ + entry.Unlock() } return nil diff --git a/prog/app.go b/prog/app.go index 7e08fdb83..5648b57d9 100644 --- a/prog/app.go +++ b/prog/app.go @@ -89,7 +89,7 @@ func router(collector app.Collector, controlRouter app.ControlRouter, pipeRouter return middlewares.Wrap(router) } -func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL, natsHostname string, +func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL string, storeInterval time.Duration, natsHostname string, memcacheConfig multitenant.MemcacheConfig, window time.Duration, maxTopNodes int, createTables bool) (app.Collector, error) { if collectorURL == "local" { return app.NewCollector(window), nil @@ -129,6 +129,7 @@ func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL, natsHo DynamoDBConfig: dynamoDBConfig, DynamoTable: tableName, S3Store: &s3Store, + StoreInterval: storeInterval, NatsHost: natsHostname, MemcacheClient: memcacheClient, Window: window, @@ -238,7 +239,7 @@ func appMain(flags appFlags) { } collector, err := collectorFactory( - userIDer, flags.collectorURL, flags.s3URL, flags.natsHostname, + userIDer, flags.collectorURL, flags.s3URL, flags.storeInterval, flags.natsHostname, multitenant.MemcacheConfig{ Host: flags.memcachedHostname, Timeout: flags.memcachedTimeout, diff --git a/prog/main.go b/prog/main.go index e44755228..620d07f46 100644 --- a/prog/main.go +++ b/prog/main.go @@ -165,6 +165,7 @@ type appFlags struct { collectorURL string s3URL string + storeInterval time.Duration controlRouterURL string controlRPCTimeout time.Duration pipeRouterURL string @@ -378,6 +379,7 @@ func setupFlags(flags *flags) { flag.StringVar(&flags.app.collectorURL, "app.collector", "local", "Collector to use (local, dynamodb, or file/directory)") flag.StringVar(&flags.app.s3URL, "app.collector.s3", "local", "S3 URL to use (when collector is dynamodb)") + flag.DurationVar(&flags.app.storeInterval, "app.collector.store-interval", 0, "How often to store merged incoming reports. If 0, reports are stored unmerged as they arrive.") flag.StringVar(&flags.app.controlRouterURL, "app.control.router", "local", "Control router to use (local or sqs)") flag.DurationVar(&flags.app.controlRPCTimeout, "app.control.rpctimeout", time.Minute, "Timeout for control RPC") flag.StringVar(&flags.app.pipeRouterURL, "app.pipe.router", "local", "Pipe router to use (local)") From 2629d13780e2b75549708d7e60a65ef2c45d3800 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 13 Apr 2020 14:30:36 +0000 Subject: [PATCH 4/6] Add a histogram for flush times --- app/multitenant/aws_collector.go | 51 +++++++++++++++++++------------- 1 file changed, 31 insertions(+), 20 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 9179bca86..33a82d28c 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -90,6 +90,13 @@ var ( Name: "nats_requests_total", Help: "Total count of NATS requests.", }, []string{"method", "status_code"}) + + flushDuration = instrument.NewHistogramCollectorFromOpts(prometheus.HistogramOpts{ + Namespace: "scope", + Name: "flush_duration_seconds", + Help: "Time in seconds spent flushing merged reports.", + Buckets: prometheus.DefBuckets, + }) ) func registerAWSCollectorMetrics() { @@ -102,6 +109,7 @@ func registerAWSCollectorMetrics() { prometheus.MustRegister(reportsPerUser) prometheus.MustRegister(reportSizePerUser) prometheus.MustRegister(natsRequests) + flushDuration.Register() } var registerAWSCollectorMetricsOnce sync.Once @@ -206,29 +214,32 @@ func (c *awsCollector) flushLoop() { // Range over all users (instances) that have pending reports and send to store func (c *awsCollector) flushPending(ctx context.Context) { - c.pending.Range(func(key, value interface{}) bool { - userid := key.(string) - entry := value.(*pendingEntry) + instrument.CollectedRequest(ctx, "FlushPending", flushDuration, nil, func(ctx context.Context) error { + c.pending.Range(func(key, value interface{}) bool { + userid := key.(string) + entry := value.(*pendingEntry) - entry.Lock() - rpt, count := entry.report, entry.count - entry.report, entry.count = report.MakeReport(), 0 - entry.Unlock() + entry.Lock() + rpt, count := entry.report, entry.count + entry.report, entry.count = report.MakeReport(), 0 + entry.Unlock() - if count > 0 { - buf, err := rpt.WriteBinary() - if err != nil { - log.Errorf("Could not serialise combined report: %v", err) - return true + if count > 0 { + buf, err := rpt.WriteBinary() + if err != nil { + log.Errorf("Could not serialise combined report: %v", err) + return true + } + rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) + err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf.Bytes()) + if err != nil { + log.Errorf("Could not persist combined report: %v", err) + return true + } } - rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) - err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf.Bytes()) - if err != nil { - log.Errorf("Could not persist combined report: %v", err) - return true - } - } - return true + return true + }) + return nil }) } From 9a739fda464b7b30b6fc17810b46c031fe67711b Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 13 Apr 2020 15:28:43 +0000 Subject: [PATCH 5/6] Parallelise sending merged reports to store Writes to DynamoDB and S3 can be done in parallel, which will reduce the overall flush time. --- app/multitenant/aws_collector.go | 31 +++++++++++++++++++++++++------ 1 file changed, 25 insertions(+), 6 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 33a82d28c..11b4a5ae0 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -215,6 +215,27 @@ func (c *awsCollector) flushLoop() { // Range over all users (instances) that have pending reports and send to store func (c *awsCollector) flushPending(ctx context.Context) { instrument.CollectedRequest(ctx, "FlushPending", flushDuration, nil, func(ctx context.Context) error { + type queueEntry struct { + userid string + buf []byte + } + queue := make(chan queueEntry) + const numParallel = 10 + var group sync.WaitGroup + group.Add(numParallel) + // Run n parallel goroutines fetching reports from the queue and flushing them + for i := 0; i < numParallel; i++ { + go func() { + for entry := range queue { + rowKey, colKey, reportKey := calculateReportKeys(entry.userid, time.Now()) + err := c.persistReport(ctx, entry.userid, rowKey, colKey, reportKey, entry.buf) + if err != nil { + log.Errorf("Could not persist combined report: %v", err) + } + } + group.Done() + }() + } c.pending.Range(func(key, value interface{}) bool { userid := key.(string) entry := value.(*pendingEntry) @@ -225,20 +246,18 @@ func (c *awsCollector) flushPending(ctx context.Context) { entry.Unlock() if count > 0 { + // serialise reports on one goroutine to limit CPU usage buf, err := rpt.WriteBinary() if err != nil { log.Errorf("Could not serialise combined report: %v", err) return true } - rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) - err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf.Bytes()) - if err != nil { - log.Errorf("Could not persist combined report: %v", err) - return true - } + queue <- queueEntry{userid: userid, buf: buf.Bytes()} } return true }) + close(queue) + group.Wait() return nil }) } From b1fc59819ac7396590088074c5233a0c6b886dec Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Wed, 15 Apr 2020 16:49:02 +0000 Subject: [PATCH 6/6] comment: clarify memcached error cases --- app/multitenant/aws_collector.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 11b4a5ae0..dc3364e8e 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -587,6 +587,9 @@ func (c *awsCollector) persistReport(ctx context.Context, userid, rowKey, colKey if c.cfg.MemcacheClient != nil { _, err = c.cfg.MemcacheClient.StoreReportBytes(ctx, reportKey, buf) if err != nil { + // NOTE: We don't abort here because failing to store in memcache + // doesn't actually break anything else -- it's just an + // optimization. log.Warningf("Could not store %v in memcache: %v", reportKey, err) } } @@ -666,7 +669,8 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e _, _, reportKey := calculateReportKeys(userid, time.Now()) _, err = c.cfg.MemcacheClient.StoreReportBytes(ctx, reportKey, buf) if err != nil { - log.Warningf("Could not store %v in memcache: %v", reportKey, err) + log.Warningf("Could not store shortcut %v in memcache: %v", reportKey, err) + // No point publishing on nats if cache store failed return nil } err := c.nats.Publish(userid, []byte(reportKey))