diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index 037b073f2..d84add02e 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -3,7 +3,9 @@ package multitenant import ( "bytes" "compress/gzip" + "crypto/md5" "fmt" + "io" "strconv" "time" @@ -11,6 +13,8 @@ import ( "github.com/aws/aws-sdk-go/aws" "github.com/aws/aws-sdk-go/aws/session" "github.com/aws/aws-sdk-go/service/dynamodb" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/bluele/gcache" "github.com/prometheus/client_golang/prometheus" "github.com/ugorji/go/codec" "golang.org/x/net/context" @@ -20,9 +24,11 @@ import ( ) const ( - hourField = "hour" - tsField = "ts" - reportField = "report" + hourField = "hour" + tsField = "ts" + reportField = "report" + cacheSize = (15 / 3) * 10 * 5 // (window size * report rate) * number of hosts per user * number of users + cacheExpiration = 15 * time.Second ) var ( @@ -31,10 +37,48 @@ var ( Name: "dynamo_request_duration_nanoseconds", Help: "Time spent doing DynamoDB requests.", }, []string{"method", "status_code"}) + dynamoCacheHits = prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "scope", + Name: "dynamo_cache_hits", + Help: "Reports fetches that hit local cache.", + }) + dynamoCacheMiss = prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "scope", + Name: "dynamo_cache_miss", + Help: "Reports fetches that miss local cache.", + }) + dynamoConsumedCapacity = prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "scope", + Name: "dynamo_consumed_capacity", + Help: "The capacity units consumed by operation.", + }, []string{"method"}) + dynamoValueSize = prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "scope", + Name: "dynamo_value_size_bytes", + Help: "Size of data read / written from dynamodb.", + }, []string{"method"}) + + reportSize = prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "scope", + Name: "report_size_bytes", + Help: "Compressed size of reports received.", + }) + + s3RequestDuration = prometheus.NewSummaryVec(prometheus.SummaryOpts{ + Namespace: "scope", + Name: "s3_request_duration_nanoseconds", + Help: "Time spent doing S3 requests.", + }, []string{"method", "status_code"}) ) func init() { prometheus.MustRegister(dynamoRequestDuration) + prometheus.MustRegister(dynamoCacheHits) + prometheus.MustRegister(dynamoCacheMiss) + prometheus.MustRegister(dynamoConsumedCapacity) + prometheus.MustRegister(dynamoValueSize) + prometheus.MustRegister(reportSize) + prometheus.MustRegister(s3RequestDuration) } // DynamoDBCollector is a Collector which can also CreateTables @@ -44,20 +88,26 @@ type DynamoDBCollector interface { } type dynamoDBCollector struct { - userIDer UserIDer - db *dynamodb.DynamoDB - tableName string - merger app.Merger + userIDer UserIDer + db *dynamodb.DynamoDB + s3 *s3.S3 + tableName string + bucketName string + merger app.Merger + cache gcache.Cache } // NewDynamoDBCollector the reaper of souls // https://github.com/aws/aws-sdk-go/wiki/common-examples -func NewDynamoDBCollector(config *aws.Config, userIDer UserIDer, tableName string) DynamoDBCollector { +func NewDynamoDBCollector(dynamoDBConfig, s3Config *aws.Config, userIDer UserIDer, tableName, bucketName string) DynamoDBCollector { return &dynamoDBCollector{ - db: dynamodb.New(session.New(config)), - userIDer: userIDer, - tableName: tableName, - merger: app.NewSmartMerger(), + db: dynamodb.New(session.New(dynamoDBConfig)), + s3: s3.New(session.New(s3Config)), + userIDer: userIDer, + tableName: tableName, + bucketName: bucketName, + merger: app.NewSmartMerger(), + cache: gcache.New(cacheSize).LRU().Expiration(cacheExpiration).Build(), } } @@ -90,7 +140,7 @@ func (c *dynamoDBCollector) CreateTables() error { // Don't need to specify non-key attributes in schema //{ // AttributeName: aws.String(reportField), - // AttributeType: aws.String("B"), + // AttributeType: aws.String("S"), //}, }, KeySchema: []*dynamodb.KeySchemaElement{ @@ -113,42 +163,99 @@ func (c *dynamoDBCollector) CreateTables() error { return err } -func (c *dynamoDBCollector) getRows(userid string, row int64, start, end time.Time) ([]report.Report, error) { - rowKey := fmt.Sprintf("%s-%s", userid, strconv.FormatInt(row, 10)) +func errorCode(err error) string { + if err == nil { + return "200" + } + return "500" +} + +func timeRequest(method string, metric *prometheus.SummaryVec, f func() error) error { startTime := time.Now() - resp, err := c.db.Query(&dynamodb.QueryInput{ - TableName: aws.String(c.tableName), - KeyConditions: map[string]*dynamodb.Condition{ - hourField: { - AttributeValueList: []*dynamodb.AttributeValue{ - {S: aws.String(rowKey)}, - }, - ComparisonOperator: aws.String("EQ"), - }, - tsField: { - AttributeValueList: []*dynamodb.AttributeValue{ - {N: aws.String(strconv.FormatInt(start.UnixNano(), 10))}, - {N: aws.String(strconv.FormatInt(end.UnixNano(), 10))}, - }, - ComparisonOperator: aws.String("BETWEEN"), - }, - }, - }) + err := f() duration := time.Now().Sub(startTime) + metric.WithLabelValues(method, errorCode(err)).Observe(float64(duration.Nanoseconds())) + return err +} + +// getReportKeys gets the s3 keys for reports in this range +func (c *dynamoDBCollector) getReportKeys(rowKey string, start, end time.Time) ([]string, error) { + var resp *dynamodb.QueryOutput + err := timeRequest("Query", dynamoRequestDuration, func() error { + var err error + resp, err = c.db.Query(&dynamodb.QueryInput{ + TableName: aws.String(c.tableName), + KeyConditions: map[string]*dynamodb.Condition{ + hourField: { + AttributeValueList: []*dynamodb.AttributeValue{ + {S: aws.String(rowKey)}, + }, + ComparisonOperator: aws.String("EQ"), + }, + tsField: { + AttributeValueList: []*dynamodb.AttributeValue{ + {N: aws.String(strconv.FormatInt(start.UnixNano(), 10))}, + {N: aws.String(strconv.FormatInt(end.UnixNano(), 10))}, + }, + ComparisonOperator: aws.String("BETWEEN"), + }, + }, + ReturnConsumedCapacity: aws.String(dynamodb.ReturnConsumedCapacityTotal), + }) + return err + }) + if resp.ConsumedCapacity != nil { + dynamoConsumedCapacity.WithLabelValues("Query"). + Add(float64(*resp.ConsumedCapacity.CapacityUnits)) + } if err != nil { - dynamoRequestDuration.WithLabelValues("Query", "500").Observe(float64(duration.Nanoseconds())) return nil, err } - dynamoRequestDuration.WithLabelValues("Query", "200").Observe(float64(duration.Nanoseconds())) - reports := []report.Report{} + + result := []string{} for _, item := range resp.Items { - b := item[reportField].B - if b == nil { + reportKey := item[reportField].S + if reportKey == nil { log.Errorf("Empty row!") continue } - buf := bytes.NewBuffer(b) - reader, err := gzip.NewReader(buf) + dynamoValueSize.WithLabelValues("BatchGetItem"). + Add(float64(len(*reportKey))) + result = append(result, *reportKey) + } + return result, nil +} + +func (c *dynamoDBCollector) getCached(reportKeys []string) ([]report.Report, []string) { + foundReports := []report.Report{} + missingReports := []string{} + for _, reportKey := range reportKeys { + rpt, err := c.cache.Get(reportKey) + if err == nil { + foundReports = append(foundReports, rpt.(report.Report)) + } else { + missingReports = append(missingReports, reportKey) + } + } + return foundReports, missingReports +} + +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 + }) + if err != nil { + return nil, err + } + reader, err := gzip.NewReader(resp.Body) if err != nil { log.Errorf("Error gunzipping report: %v", err) continue @@ -159,10 +266,33 @@ func (c *dynamoDBCollector) getRows(userid string, row int64, start, end time.Ti continue } reports = append(reports, rep) + c.cache.Set(reportKey, rep) } return reports, 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) + if err != nil { + return nil, err + } + + cachedReports, missing := c.getCached(reportKeys) + dynamoCacheHits.Add(float64(len(cachedReports))) + dynamoCacheMiss.Add(float64(len(missing))) + if len(missing) == 0 { + return cachedReports, nil + } + + fetchedReport, err := c.getNonCached(missing) + if err != nil { + return nil, err + } + + return append(cachedReports, fetchedReport...), nil +} + func (c *dynamoDBCollector) Report(ctx context.Context) (report.Report, error) { var ( now = time.Now() @@ -177,19 +307,19 @@ func (c *dynamoDBCollector) Report(ctx context.Context) (report.Report, error) { // Queries will only every span 2 rows max. if rowStart != rowEnd { - reports1, err := c.getRows(userid, rowStart, start, now) + reports1, err := c.getReports(userid, rowStart, start, now) if err != nil { return report.MakeReport(), err } - reports2, err := c.getRows(userid, rowEnd, start, now) + reports2, err := c.getReports(userid, rowEnd, start, now) if err != nil { return report.MakeReport(), err } reports = append(reports1, reports2...) } else { - if reports, err = c.getRows(userid, rowEnd, start, now); err != nil { + if reports, err = c.getReports(userid, rowEnd, start, now); err != nil { return report.MakeReport(), err } } @@ -203,37 +333,69 @@ func (c *dynamoDBCollector) Add(ctx context.Context, rep report.Report) error { return err } + // first, encode the report into a buffer and record its size var buf bytes.Buffer - writer := gzip.NewWriter(&buf) + writer, err := gzip.NewWriterLevel(&buf, gzip.BestCompression) + if err != nil { + return err + } if err := codec.NewEncoder(writer, &codec.MsgpackHandle{}).Encode(&rep); err != nil { return err } writer.Close() + reportSize.Add(float64(buf.Len())) + // second, put the report on s3 now := time.Now() rowKey := fmt.Sprintf("%s-%s", userid, strconv.FormatInt(now.UnixNano()/time.Hour.Nanoseconds(), 10)) - startTime := time.Now() - _, err = c.db.PutItem(&dynamodb.PutItemInput{ - TableName: aws.String(c.tableName), - Item: map[string]*dynamodb.AttributeValue{ - hourField: { - S: aws.String(rowKey), - }, - tsField: { - N: aws.String(strconv.FormatInt(now.UnixNano(), 10)), - }, - reportField: { - B: buf.Bytes(), - }, - }, - }) - duration := time.Now().Sub(startTime) - if err != nil { - dynamoRequestDuration.WithLabelValues("PutItem", "500").Observe(float64(duration.Nanoseconds())) + colKey := strconv.FormatInt(now.UnixNano(), 10) + rowKeyHash := md5.New() + if _, err := io.WriteString(rowKeyHash, rowKey); err != nil { return err } - dynamoRequestDuration.WithLabelValues("PutItem", "200").Observe(float64(duration.Nanoseconds())) - return nil + s3Key := fmt.Sprintf("%x/%s", rowKeyHash.Sum(nil), colKey) + err = timeRequest("Put", s3RequestDuration, func() error { + var err error + _, err = c.s3.PutObject(&s3.PutObjectInput{ + Body: bytes.NewReader(buf.Bytes()), + Bucket: aws.String(c.bucketName), + Key: aws.String(s3Key), + }) + return err + }) + if err != nil { + return err + } + + // third, put the key in dynamodb + dynamoValueSize.WithLabelValues("PutItem"). + Add(float64(len(s3Key))) + + var resp *dynamodb.PutItemOutput + err = timeRequest("PutItem", dynamoRequestDuration, func() error { + var err error + resp, err = c.db.PutItem(&dynamodb.PutItemInput{ + TableName: aws.String(c.tableName), + Item: map[string]*dynamodb.AttributeValue{ + hourField: { + S: aws.String(rowKey), + }, + tsField: { + N: aws.String(colKey), + }, + reportField: { + S: aws.String(s3Key), + }, + }, + ReturnConsumedCapacity: aws.String(dynamodb.ReturnConsumedCapacityTotal), + }) + return err + }) + if resp.ConsumedCapacity != nil { + dynamoConsumedCapacity.WithLabelValues("PutItem"). + Add(float64(*resp.ConsumedCapacity.CapacityUnits)) + } + return err } func (c *dynamoDBCollector) WaitOn(context.Context, chan struct{}) {} diff --git a/app/router.go b/app/router.go index eb0f1a6c1..036912cbc 100644 --- a/app/router.go +++ b/app/router.go @@ -160,7 +160,11 @@ func RegisterReportPostHandler(a Adder, router *mux.Router) { float32(compressedSize)/float32(uncompressedSize)*100, ) - a.Add(ctx, rpt) + if err := a.Add(ctx, rpt); err != nil { + log.Errorf("Error Adding report: %v", err) + http.Error(w, err.Error(), http.StatusBadRequest) + return + } w.WriteHeader(http.StatusOK) })) } diff --git a/prog/app.go b/prog/app.go index 22df1e1e0..c3f9d7cc8 100644 --- a/prog/app.go +++ b/prog/app.go @@ -61,7 +61,10 @@ func router(collector app.Collector, controlRouter app.ControlRouter, pipeRouter return instrument.Wrap(router) } -func awsConfigFromURL(url *url.URL) *aws.Config { +func awsConfigFromURL(url *url.URL) (*aws.Config, error) { + if url.User == nil { + return nil, fmt.Errorf("Must specify username & password in URL") + } password, _ := url.User.Password() creds := credentials.NewStaticCredentials(url.User.Username(), password, "") config := aws.NewConfig().WithCredentials(creds) @@ -70,10 +73,10 @@ func awsConfigFromURL(url *url.URL) *aws.Config { } else { config = config.WithRegion(url.Host) } - return config + return config, nil } -func collectorFactory(userIDer multitenant.UserIDer, collectorURL string, window time.Duration, createTables bool) (app.Collector, error) { +func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL string, window time.Duration, createTables bool) (app.Collector, error) { if collectorURL == "local" { return app.NewCollector(window), nil } @@ -84,8 +87,22 @@ func collectorFactory(userIDer multitenant.UserIDer, collectorURL string, window } if parsed.Scheme == "dynamodb" { + s3, err := url.Parse(s3URL) + if err != nil { + return nil, fmt.Errorf("Valid URL for s3 required: %v", err) + } + dynamoDBConfig, err := awsConfigFromURL(parsed) + if err != nil { + return nil, err + } + s3Config, err := awsConfigFromURL(s3) + if err != nil { + return nil, err + } tableName := strings.TrimPrefix(parsed.Path, "/") - dynamoCollector := multitenant.NewDynamoDBCollector(awsConfigFromURL(parsed), userIDer, tableName) + bucketName := strings.TrimPrefix(s3.Path, "/") + dynamoCollector := multitenant.NewDynamoDBCollector( + dynamoDBConfig, s3Config, userIDer, tableName, bucketName) if createTables { if err := dynamoCollector.CreateTables(); err != nil { return nil, err @@ -109,7 +126,11 @@ func controlRouterFactory(userIDer multitenant.UserIDer, controlRouterURL string if parsed.Scheme == "sqs" { prefix := strings.TrimPrefix(parsed.Path, "/") - return multitenant.NewSQSControlRouter(awsConfigFromURL(parsed), userIDer, prefix), nil + sqsConfig, err := awsConfigFromURL(parsed) + if err != nil { + return nil, err + } + return multitenant.NewSQSControlRouter(sqsConfig, userIDer, prefix), nil } return nil, fmt.Errorf("Invalid control router '%s'", controlRouterURL) @@ -151,7 +172,7 @@ func appMain(flags appFlags) { userIDer = multitenant.UserIDHeader(flags.userIDHeader) } - collector, err := collectorFactory(userIDer, flags.collectorURL, flags.window, flags.awsCreateTables) + collector, err := collectorFactory(userIDer, flags.collectorURL, flags.s3URL, flags.window, flags.awsCreateTables) if err != nil { log.Fatalf("Error creating collector: %v", err) return diff --git a/prog/main.go b/prog/main.go index 382ec6401..285e33912 100644 --- a/prog/main.go +++ b/prog/main.go @@ -97,6 +97,7 @@ type appFlags struct { dockerEndpoint string collectorURL string + s3URL string controlRouterURL string pipeRouterURL string userIDHeader string @@ -172,6 +173,7 @@ func main() { flag.StringVar(&flags.app.dockerEndpoint, "app.docker", app.DefaultDockerEndpoint, "Location of docker endpoint (to lookup container ID)") flag.StringVar(&flags.app.collectorURL, "app.collector", "local", "Collector to use (local of dynamodb)") + flag.StringVar(&flags.app.s3URL, "app.collector.s3", "local", "S3 URL to use (when collector is dynamodb)") flag.StringVar(&flags.app.controlRouterURL, "app.control.router", "local", "Control router to use (local or sqs)") flag.StringVar(&flags.app.pipeRouterURL, "app.pipe.router", "local", "Pipe router to use (local)") flag.StringVar(&flags.app.userIDHeader, "app.userid.header", "", "HTTP header to use as userid")