Merge pull request #3780 from weaveworks/merge-in-collector

Multitenant: merge incoming reports in collector
This commit is contained in:
Bryan Boreham
2020-04-15 19:08:39 +01:00
committed by GitHub
5 changed files with 177 additions and 63 deletions
+7
View File
@@ -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) {
+161 -58
View File
@@ -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
@@ -123,17 +131,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
@@ -160,6 +178,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 {
@@ -169,14 +190,82 @@ 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
}
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) {
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)
entry.Lock()
rpt, count := entry.report, entry.count
entry.report, entry.count = report.MakeReport(), 0
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
}
queue <- queueEntry{userid: userid, buf: buf.Bytes()}
}
return true
})
close(queue)
group.Wait()
return nil
})
}
// 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
@@ -481,13 +570,49 @@ 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 {
// 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)
}
}
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 +662,41 @@ 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 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 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))
natsRequests.WithLabelValues("Publish", instrument.ErrorCode(err)).Add(1)
if err != nil {
log.Errorf("Error sending shortcut report: %v", err)
}
}
return nil
}
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
}
reportSizeHistogram.Observe(float64(reportSize))
reportSizePerUser.WithLabelValues(userid).Add(float64(reportSize))
reportsPerUser.WithLabelValues(userid).Inc()
}
// 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)
} 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
+3 -2
View File
@@ -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()
}
+4 -3
View File
@@ -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,
@@ -259,9 +260,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 {
+2
View File
@@ -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)")