Put reports in S3; add in process caching (#1545)

* Add in-process caching to dynamodb collector

* Add metrics for dynamodb consumed capacity and report size

* Log and return errors during report collection

* Increase compression to the max

* Put reports in S3 and just use DynamoDB as an index.

* Review feedback
This commit is contained in:
Tom Wilkie
2016-05-31 15:40:15 +01:00
parent 378e2bd24e
commit 12f281654d
4 changed files with 260 additions and 71 deletions

View File

@@ -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{}) {}

View File

@@ -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)
}))
}

View File

@@ -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

View File

@@ -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")