From 3d12a2a76c72bd7fd7ab69f40121372e0014aa1c Mon Sep 17 00:00:00 2001 From: Jonathan Lange Date: Tue, 7 Jun 2016 18:48:01 +0100 Subject: [PATCH 1/3] Extract function for getting single report --- app/multitenant/dynamo_collector.go | 47 +++++++++++++++++------------ 1 file changed, 28 insertions(+), 19 deletions(-) diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index d84add02e..d84c4818b 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -243,34 +243,43 @@ func (c *dynamoDBCollector) getCached(reportKeys []string) ([]report.Report, []s func (c *dynamoDBCollector) getNonCached(reportKeys []string) ([]report.Report, error) { reports := []report.Report{} for _, reportKey := range reportKeys { - var resp *s3.GetObjectOutput - err := timeRequest("Get", s3RequestDuration, func() error { - var err error - resp, err = c.s3.GetObject(&s3.GetObjectInput{ - Bucket: aws.String(c.bucketName), - Key: aws.String(reportKey), - }) - return err - }) + rep, err := c.getNonCachedReport(reportKey) if err != nil { return nil, err } - reader, err := gzip.NewReader(resp.Body) - if err != nil { - log.Errorf("Error gunzipping report: %v", err) - continue - } - rep := report.MakeReport() - if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { - log.Errorf("Failed to decode report: %v", err) - continue - } reports = append(reports, rep) c.cache.Set(reportKey, rep) } return reports, nil } +func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (report.Report, error) { + var resp *s3.GetObjectOutput + err := timeRequest("Get", s3RequestDuration, func() error { + var err error + resp, err = c.s3.GetObject(&s3.GetObjectInput{ + Bucket: aws.String(c.bucketName), + Key: aws.String(reportKey), + }) + return err + }) + // XXX: Not sure why this one is returned out when the rest are logged. + if err != nil { + return report.Report{}, err + } + reader, err := gzip.NewReader(resp.Body) + if err != nil { + log.Errorf("Error gunzipping report: %v", err) + return report.Report{}, nil + } + rep := report.MakeReport() + if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { + log.Errorf("Failed to decode report: %v", err) + return report.Report{}, nil + } + return rep, nil +} + func (c *dynamoDBCollector) getReports(userid string, row int64, start, end time.Time) ([]report.Report, error) { rowKey := fmt.Sprintf("%s-%s", userid, strconv.FormatInt(row, 10)) reportKeys, err := c.getReportKeys(rowKey, start, end) From 0907cdfa0dd4be47e75ba90566e3006eb14d9842 Mon Sep 17 00:00:00 2001 From: Jonathan Lange Date: Tue, 7 Jun 2016 19:00:14 +0100 Subject: [PATCH 2/3] Fail fast on error fetching non-cached reports --- app/multitenant/dynamo_collector.go | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index d84c4818b..81190b9f9 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -253,6 +253,7 @@ func (c *dynamoDBCollector) getNonCached(reportKeys []string) ([]report.Report, return reports, nil } +// Fetch a single report from S3. func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (report.Report, error) { var resp *s3.GetObjectOutput err := timeRequest("Get", s3RequestDuration, func() error { @@ -263,19 +264,16 @@ func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (report.Report, }) return err }) - // XXX: Not sure why this one is returned out when the rest are logged. if err != nil { return report.Report{}, err } reader, err := gzip.NewReader(resp.Body) if err != nil { - log.Errorf("Error gunzipping report: %v", err) - return report.Report{}, nil + return report.Report{}, err } rep := report.MakeReport() if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { - log.Errorf("Failed to decode report: %v", err) - return report.Report{}, nil + return report.Report{}, err } return rep, nil } From 48fc985a3e165eacb380ddcb63bb373617a4e08f Mon Sep 17 00:00:00 2001 From: Jonathan Lange Date: Tue, 7 Jun 2016 19:23:45 +0100 Subject: [PATCH 3/3] Get non-cached reports in parallel --- app/multitenant/dynamo_collector.go | 39 +++++++++++++++++++++-------- 1 file changed, 28 insertions(+), 11 deletions(-) diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index 81190b9f9..d9d88b03c 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -240,21 +240,38 @@ func (c *dynamoDBCollector) getCached(reportKeys []string) ([]report.Report, []s return foundReports, missingReports } +// Fetch multiple reports in parallel from S3. func (c *dynamoDBCollector) getNonCached(reportKeys []string) ([]report.Report, error) { - reports := []report.Report{} + type result struct { + key string + report *report.Report + err error + } + + ch := make(chan result, len(reportKeys)) + for _, reportKey := range reportKeys { - rep, err := c.getNonCachedReport(reportKey) - if err != nil { - return nil, err + go func(reportKey string) { + r := result{key: reportKey} + r.report, r.err = c.getNonCachedReport(reportKey) + ch <- r + }(reportKey) + } + + reports := []report.Report{} + for range reportKeys { + r := <-ch + if r.err != nil { + return nil, r.err } - reports = append(reports, rep) - c.cache.Set(reportKey, rep) + reports = append(reports, *r.report) + c.cache.Set(r.key, *r.report) } return reports, nil } // Fetch a single report from S3. -func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (report.Report, error) { +func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (*report.Report, error) { var resp *s3.GetObjectOutput err := timeRequest("Get", s3RequestDuration, func() error { var err error @@ -265,17 +282,17 @@ func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (report.Report, return err }) if err != nil { - return report.Report{}, err + return nil, err } reader, err := gzip.NewReader(resp.Body) if err != nil { - return report.Report{}, err + return nil, err } rep := report.MakeReport() if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { - return report.Report{}, err + return nil, err } - return rep, nil + return &rep, nil } func (c *dynamoDBCollector) getReports(userid string, row int64, start, end time.Time) ([]report.Report, error) {