Rate-limit reading proc files

Use a reader in the background, dynamically rate-limited, reading the required
files in a loop
This commit is contained in:
Alfonso Acosta
2016-02-04 12:33:04 +00:00
parent 2695a1a2e7
commit f922ea19c8
4 changed files with 118 additions and 4 deletions

View File

@@ -0,0 +1,110 @@
package procspy
import (
"bytes"
"fmt"
"log"
"math"
"sync"
"time"
"github.com/weaveworks/scope/probe/process"
)
const (
initialRateLimit = 100 * time.Millisecond // read 10 namespaces per second
maxRateLimit = 250 * time.Millisecond // lead at least 4 namespaces per second
targetWalkTime = 10 * time.Second
maxRateLimitF = float64(maxRateLimit)
targetWalkTimeF = float64(targetWalkTime)
)
type backgroundReader struct {
walker process.Walker
mtx sync.Mutex
walkingBuf *bytes.Buffer
readyBuf *bytes.Buffer
readySockets map[uint64]*Proc
}
// HACK: Pretty ugly singleton interface (particularly the part of passing
// the walker to StartBackgroundReader() and ignoring it in in Connections() )
// experimenting with this for now.
var singleton *backgroundReader
func getBackgroundReader() (*backgroundReader, error) {
var err error
if singleton == nil {
err = fmt.Errorf("background reader hasn't yet been started")
}
return singleton, err
}
// StartBackgroundReader starts a ratelimited background goroutine to
// read the expensive files from proc.
func StartBackgroundReader(walker process.Walker) {
if singleton != nil {
return
}
singleton = &backgroundReader{
walker: walker,
walkingBuf: bytes.NewBuffer(make([]byte, 0, 5000)),
readyBuf: bytes.NewBuffer(make([]byte, 0, 5000)),
}
go singleton.loop()
}
func (br *backgroundReader) loop() {
rateLimit := initialRateLimit
namespaceTicker := time.Tick(rateLimit)
for {
start := time.Now()
sockets, err := walkProcPid(br.walkingBuf, br.walker, namespaceTicker)
if err != nil {
log.Printf("background reader: error walking /proc: %s\n", err)
continue
}
walkTime := time.Now().Sub(start)
walkTimeF := float64(walkTime)
log.Printf("debug: background reader: full pass took %s\n", walkTime)
if walkTimeF/targetWalkTimeF > 1.5 {
log.Printf(
"warn: background reader: full pass took %s: 50%% more than expected (%s)\n",
walkTime,
targetWalkTime,
)
}
// Adjust rate limit to more-accurately meet the target walk time in next iteration
scaledRateLimit := targetWalkTimeF / walkTimeF * float64(rateLimit)
rateLimit = time.Duration(math.Min(scaledRateLimit, maxRateLimitF))
log.Printf("debug: background reader: new rate limit %s\n", rateLimit)
namespaceTicker = time.Tick(rateLimit)
// Swap buffers
br.mtx.Lock()
br.readyBuf, br.walkingBuf = br.walkingBuf, br.readyBuf
br.readySockets = sockets
br.mtx.Unlock()
br.walkingBuf.Reset()
// Sleep during spare time
time.Sleep(targetWalkTime - walkTime)
}
}
func (br *backgroundReader) getWalkedProcPid(buf *bytes.Buffer) map[uint64]*Proc {
br.mtx.Lock()
defer br.mtx.Unlock()
reader := bytes.NewReader(br.readyBuf.Bytes())
buf.ReadFrom(reader)
return br.readySockets
}

View File

@@ -7,6 +7,7 @@ import (
"path/filepath"
"strconv"
"syscall"
"time"
log "github.com/Sirupsen/logrus"
"github.com/armon/go-metrics"
@@ -143,7 +144,7 @@ func walkNamespacePid(buf *bytes.Buffer, sockets map[uint64]*Proc, namespaceProc
// /proc/PID/net/tcp{,6} for each namespace and sees if the ./fd/* files of each
// process in that namespace are symlinks to sockets. Returns a map from socket
// ID (inode) to PID.
func walkProcPid(buf *bytes.Buffer, walker process.Walker) (map[uint64]*Proc, error) {
func walkProcPid(buf *bytes.Buffer, walker process.Walker, namespaceTicker <-chan time.Time) (map[uint64]*Proc, error) {
var (
sockets = map[uint64]*Proc{} // map socket inode -> process
namespaces = map[uint64][]*process.Process{} // map network namespace id -> processes
@@ -171,6 +172,7 @@ func walkProcPid(buf *bytes.Buffer, walker process.Walker) (map[uint64]*Proc, er
})
for _, procs := range namespaces {
<-namespaceTicker
walkNamespacePid(buf, sockets, procs)
}

View File

@@ -33,17 +33,18 @@ func (c *pnConnIter) Next() *Connection {
}
// cbConnections sets Connections()
var cbConnections = func(processes bool, walker process.Walker) (ConnIter, error) {
var cbConnections = func(processes bool, _ process.Walker) (ConnIter, error) {
// buffer for contents of /proc/<pid>/net/tcp
buf := bufPool.Get().(*bytes.Buffer)
buf.Reset()
var procs map[uint64]*Proc
if processes {
var err error
if procs, err = walkProcPid(buf, walker); err != nil {
br, err := getBackgroundReader()
if err != nil {
return nil, err
}
procs = br.getWalkedProcPid(buf)
}
if buf.Len() == 0 {

View File

@@ -49,6 +49,7 @@ var SpyDuration = prometheus.NewSummaryVec(
// is stored in the Endpoint topology. It optionally enriches that topology
// with process (PID) information.
func NewReporter(hostID, hostName string, includeProcesses bool, useConntrack bool, procWalker process.Walker) *Reporter {
procspy.StartBackgroundReader(procWalker)
return &Reporter{
hostID: hostID,
hostName: hostName,