This commit is contained in:
gary rong 2018-01-29 10:00:36 +00:00 committed by GitHub
commit 330d9ce0d9
7 changed files with 148 additions and 72 deletions

View file

@ -30,24 +30,31 @@ import (
"github.com/syndtr/goleveldb/leveldb/iterator" "github.com/syndtr/goleveldb/leveldb/iterator"
"github.com/syndtr/goleveldb/leveldb/opt" "github.com/syndtr/goleveldb/leveldb/opt"
"fmt"
gometrics "github.com/rcrowley/go-metrics" gometrics "github.com/rcrowley/go-metrics"
) )
var OpenFileLimit = 64 const (
writeDelayNThreshold = 200
writeDelayThreshold = 350 * time.Millisecond
writeDelayWarningThrottler = 1 * time.Minute
)
type LDBDatabase struct { type LDBDatabase struct {
fn string // filename for reporting fn string // filename for reporting
db *leveldb.DB // LevelDB instance db *leveldb.DB // LevelDB instance
getTimer gometrics.Timer // Timer for measuring the database get request counts and latencies getTimer gometrics.Timer // Timer for measuring the database get request counts and latencies
putTimer gometrics.Timer // Timer for measuring the database put request counts and latencies putTimer gometrics.Timer // Timer for measuring the database put request counts and latencies
delTimer gometrics.Timer // Timer for measuring the database delete request counts and latencies delTimer gometrics.Timer // Timer for measuring the database delete request counts and latencies
missMeter gometrics.Meter // Meter for measuring the missed database get requests missMeter gometrics.Meter // Meter for measuring the missed database get requests
readMeter gometrics.Meter // Meter for measuring the database get request data usage readMeter gometrics.Meter // Meter for measuring the database get request data usage
writeMeter gometrics.Meter // Meter for measuring the database put request data usage writeMeter gometrics.Meter // Meter for measuring the database put request data usage
compTimeMeter gometrics.Meter // Meter for measuring the total time spent in database compaction writeDelayNMeter gometrics.Meter // Meter for measuring the write delay number due to database compaction
compReadMeter gometrics.Meter // Meter for measuring the data read during compaction writeDelayMeter gometrics.Meter // Meter for measuring the write delay duration due to database compaction
compWriteMeter gometrics.Meter // Meter for measuring the data written during compaction compTimeMeter gometrics.Meter // Meter for measuring the total time spent in database compaction
compReadMeter gometrics.Meter // Meter for measuring the data read during compaction
compWriteMeter gometrics.Meter // Meter for measuring the data written during compaction
quitLock sync.Mutex // Mutex protecting the quit channel access quitLock sync.Mutex // Mutex protecting the quit channel access
quitChan chan chan error // Quit channel to stop the metrics collection before closing the database quitChan chan chan error // Quit channel to stop the metrics collection before closing the database
@ -175,20 +182,23 @@ func (db *LDBDatabase) LDB() *leveldb.DB {
// Meter configures the database metrics collectors and // Meter configures the database metrics collectors and
func (db *LDBDatabase) Meter(prefix string) { func (db *LDBDatabase) Meter(prefix string) {
// Short circuit metering if the metrics system is disabled if metrics.Enabled {
if !metrics.Enabled { // Initialize all database related metrics collector at the requested prefix
return // if metric is enable.
db.getTimer = metrics.NewTimer(prefix + "user/gets")
db.putTimer = metrics.NewTimer(prefix + "user/puts")
db.delTimer = metrics.NewTimer(prefix + "user/dels")
db.missMeter = metrics.NewMeter(prefix + "user/misses")
db.readMeter = metrics.NewMeter(prefix + "user/reads")
db.writeMeter = metrics.NewMeter(prefix + "user/writes")
db.compTimeMeter = metrics.NewMeter(prefix + "compact/time")
db.compReadMeter = metrics.NewMeter(prefix + "compact/input")
db.compWriteMeter = metrics.NewMeter(prefix + "compact/output")
} }
// Initialize all the metrics collector at the requested prefix // Initialize essential write delay metrics no matter metric is enable or not.
db.getTimer = metrics.NewTimer(prefix + "user/gets") db.writeDelayMeter = metrics.NewMeter(prefix + "compact/writedelay/duration")
db.putTimer = metrics.NewTimer(prefix + "user/puts") db.writeDelayNMeter = metrics.NewMeter(prefix + "compact/writedelay/counter")
db.delTimer = metrics.NewTimer(prefix + "user/dels")
db.missMeter = metrics.NewMeter(prefix + "user/misses")
db.readMeter = metrics.NewMeter(prefix + "user/reads")
db.writeMeter = metrics.NewMeter(prefix + "user/writes")
db.compTimeMeter = metrics.NewMeter(prefix + "compact/time")
db.compReadMeter = metrics.NewMeter(prefix + "compact/input")
db.compWriteMeter = metrics.NewMeter(prefix + "compact/output")
// Create a quit channel for the periodic collector and run it // Create a quit channel for the periodic collector and run it
db.quitLock.Lock() db.quitLock.Lock()
@ -210,58 +220,114 @@ func (db *LDBDatabase) Meter(prefix string) {
// 2 | 523 | 1000.37159 | 7.26059 | 66.86342 | 66.77884 // 2 | 523 | 1000.37159 | 7.26059 | 66.86342 | 66.77884
// 3 | 570 | 1113.18458 | 0.00000 | 0.00000 | 0.00000 // 3 | 570 | 1113.18458 | 0.00000 | 0.00000 | 0.00000
func (db *LDBDatabase) meter(refresh time.Duration) { func (db *LDBDatabase) meter(refresh time.Duration) {
var (
lastWriteDelay time.Time
lastWriteDelayN time.Time
)
// Create the counters to store current and previous values // Create the counters to store current and previous values
counters := make([][]float64, 2) counters := make([][]float64, 2)
for i := 0; i < 2; i++ { for i := 0; i < 2; i++ {
counters[i] = make([]float64, 3) counters[i] = make([]float64, 5)
} }
// Iterate ad infinitum and collect the stats // Iterate ad infinitum and collect the stats
for i := 1; ; i++ { for i := 1; ; i++ {
// Retrieve the database stats if metrics.Enabled {
stats, err := db.db.GetProperty("leveldb.stats") // Retrieve the database stats
if err != nil { stats, err := db.db.GetProperty("leveldb.stats")
db.log.Error("Failed to read database stats", "err", err) if err != nil {
return db.log.Error("Failed to read database stats", "err", err)
} return
// Find the compaction table, skip the header }
lines := strings.Split(stats, "\n") // Find the compaction table, skip the header
for len(lines) > 0 && strings.TrimSpace(lines[0]) != "Compactions" { lines := strings.Split(stats, "\n")
lines = lines[1:] for len(lines) > 0 && strings.TrimSpace(lines[0]) != "Compactions" {
} lines = lines[1:]
if len(lines) <= 3 { }
db.log.Error("Compaction table not found") if len(lines) <= 3 {
return db.log.Error("Compaction table not found")
} return
lines = lines[3:] }
lines = lines[3:]
// Iterate over all the table rows, and accumulate the entries // Iterate over all the table rows, and accumulate the entries
for j := 0; j < len(counters[i%2]); j++ { for j := 0; j < len(counters[i%2]); j++ {
counters[i%2][j] = 0 counters[i%2][j] = 0
}
for _, line := range lines {
parts := strings.Split(line, "|")
if len(parts) != 6 {
break
} }
for idx, counter := range parts[3:] { for _, line := range lines {
value, err := strconv.ParseFloat(strings.TrimSpace(counter), 64) parts := strings.Split(line, "|")
if err != nil { if len(parts) != 6 {
db.log.Error("Compaction entry parsing failed", "err", err) break
return
} }
counters[i%2][idx] += value for idx, counter := range parts[3:] {
value, err := strconv.ParseFloat(strings.TrimSpace(counter), 64)
if err != nil {
db.log.Error("Compaction entry parsing failed", "err", err)
return
}
counters[i%2][idx] += value
}
}
// Update all the requested meters
if db.compTimeMeter != nil {
db.compTimeMeter.Mark(int64((counters[i%2][0] - counters[(i-1)%2][0]) * 1000 * 1000 * 1000))
}
if db.compReadMeter != nil {
db.compReadMeter.Mark(int64((counters[i%2][1] - counters[(i-1)%2][1]) * 1024 * 1024))
}
if db.compWriteMeter != nil {
db.compWriteMeter.Mark(int64((counters[i%2][2] - counters[(i-1)%2][2]) * 1024 * 1024))
} }
} }
// Update all the requested meters
if db.compTimeMeter != nil { // Gather write delay statistic
db.compTimeMeter.Mark(int64((counters[i%2][0] - counters[(i-1)%2][0]) * 1000 * 1000 * 1000)) writeDelay, err := db.db.GetProperty("leveldb.writedelay")
if err != nil {
db.log.Error("Failed to read database write delay statistic", "err", err)
return
} }
if db.compReadMeter != nil { var (
db.compReadMeter.Mark(int64((counters[i%2][1] - counters[(i-1)%2][1]) * 1024 * 1024)) delayN int32
delayDuration string
duration time.Duration
)
if n, err := fmt.Sscanf(writeDelay, "DelayN:%d Delay:%s", &delayN, &delayDuration); n != 2 || err != nil {
db.log.Error("Write delay statistic not found")
return
} }
if db.compWriteMeter != nil { duration, err = time.ParseDuration(delayDuration)
db.compWriteMeter.Mark(int64((counters[i%2][2] - counters[(i-1)%2][2]) * 1024 * 1024)) if err != nil {
db.log.Error("Failed to parse delay duration", "err", err)
return
} }
counters[i%2][3], counters[i%2][4] = float64(delayN), float64(duration.Nanoseconds())
if db.writeDelayNMeter != nil {
db.writeDelayNMeter.Mark(int64((counters[i%2][3] - counters[(i-1)%2][3])))
// If the write delay number been collected in the last minute exceeds the predefined threshold,
// print a warning log here.
// If a warning that db performance is laggy has been displayed,
// any subsequent warnings will be withhold for 1 minute to don't overwhelm the user.
if int(db.writeDelayNMeter.Rate1()) > writeDelayNThreshold &&
time.Now().After(lastWriteDelayN.Add(writeDelayWarningThrottler)) {
db.log.Warn("Write delay number exceeds the threshold (200) in the last minute")
lastWriteDelayN = time.Now()
}
}
if db.writeDelayMeter != nil {
db.writeDelayMeter.Mark(int64((counters[i%2][4] - counters[(i-1)%2][4])))
// If the write delay duration been collected in the last minute exceeds the predefined threshold,
// print a warning log here.
// If a warning that db performance is laggy has been displayed,
// any subsequent warnings will be withhold for 1 minute to don't overwhelm the user.
if int64(db.writeDelayMeter.Rate1()) > writeDelayThreshold.Nanoseconds() &&
time.Now().After(lastWriteDelay.Add(writeDelayWarningThrottler)) {
db.log.Warn("Write delay duration exceeds the threshold (0.35sec) in the last minute")
lastWriteDelay = time.Now()
}
}
// Sleep a bit, then repeat the stats collection // Sleep a bit, then repeat the stats collection
select { select {
case errc := <-db.quitChan: case errc := <-db.quitChan:

View file

@ -32,6 +32,11 @@ type DB struct {
// Need 64-bit alignment. // Need 64-bit alignment.
seq uint64 seq uint64
// Stats. Need 64-bit alignment.
cWriteDelay int64 // The cumulative duration of write delays
cWriteDelayN int32 // The cumulative number of write delays
aliveSnaps, aliveIters int32
// Session. // Session.
s *session s *session
@ -49,9 +54,6 @@ type DB struct {
snapsMu sync.Mutex snapsMu sync.Mutex
snapsList *list.List snapsList *list.List
// Stats.
aliveSnaps, aliveIters int32
// Write. // Write.
batchPool sync.Pool batchPool sync.Pool
writeMergeC chan writeMerge writeMergeC chan writeMerge
@ -321,7 +323,7 @@ func recoverTable(s *session, o *opt.Options) error {
} }
} }
err = iter.Error() err = iter.Error()
if err != nil { if err != nil && !errors.IsCorrupted(err) {
return return
} }
err = tw.Close() err = tw.Close()
@ -392,7 +394,7 @@ func recoverTable(s *session, o *opt.Options) error {
} }
imax = append(imax[:0], key...) imax = append(imax[:0], key...)
} }
if err := iter.Error(); err != nil { if err := iter.Error(); err != nil && !errors.IsCorrupted(err) {
iter.Release() iter.Release()
return err return err
} }
@ -904,6 +906,8 @@ func (db *DB) GetSnapshot() (*Snapshot, error) {
// Returns the number of files at level 'n'. // Returns the number of files at level 'n'.
// leveldb.stats // leveldb.stats
// Returns statistics of the underlying DB. // Returns statistics of the underlying DB.
// leveldb.writedelay
// Returns cumulative write delay caused by compaction.
// leveldb.sstables // leveldb.sstables
// Returns sstables list for each level. // Returns sstables list for each level.
// leveldb.blockpool // leveldb.blockpool
@ -955,6 +959,9 @@ func (db *DB) GetProperty(name string) (value string, err error) {
level, len(tables), float64(tables.size())/1048576.0, duration.Seconds(), level, len(tables), float64(tables.size())/1048576.0, duration.Seconds(),
float64(read)/1048576.0, float64(write)/1048576.0) float64(read)/1048576.0, float64(write)/1048576.0)
} }
case p == "writedelay":
writeDelayN, writeDelay := atomic.LoadInt32(&db.cWriteDelayN), time.Duration(atomic.LoadInt64(&db.cWriteDelay))
value = fmt.Sprintf("DelayN:%d Delay:%s", writeDelayN, writeDelay)
case p == "sstables": case p == "sstables":
for level, tables := range v.levels { for level, tables := range v.levels {
value += fmt.Sprintf("--- level %d ---\n", level) value += fmt.Sprintf("--- level %d ---\n", level)

View file

@ -7,6 +7,7 @@
package leveldb package leveldb
import ( import (
"sync/atomic"
"time" "time"
"github.com/syndtr/goleveldb/leveldb/memdb" "github.com/syndtr/goleveldb/leveldb/memdb"
@ -117,6 +118,8 @@ func (db *DB) flush(n int) (mdb *memDB, mdbFree int, err error) {
db.writeDelayN++ db.writeDelayN++
} else if db.writeDelayN > 0 { } else if db.writeDelayN > 0 {
db.logf("db@write was delayed N·%d T·%v", db.writeDelayN, db.writeDelay) db.logf("db@write was delayed N·%d T·%v", db.writeDelayN, db.writeDelay)
atomic.AddInt32(&db.cWriteDelayN, int32(db.writeDelayN))
atomic.AddInt64(&db.cWriteDelay, int64(db.writeDelay))
db.writeDelay = 0 db.writeDelay = 0
db.writeDelayN = 0 db.writeDelayN = 0
} }

View file

@ -88,7 +88,7 @@ type Iterator interface {
// its contents may change on the next call to any 'seeks method'. // its contents may change on the next call to any 'seeks method'.
Key() []byte Key() []byte
// Value returns the key of the current key/value pair, or nil if done. // Value returns the value of the current key/value pair, or nil if done.
// The caller should not modify the contents of the returned slice, and // The caller should not modify the contents of the returned slice, and
// its contents may change on the next call to any 'seeks method'. // its contents may change on the next call to any 'seeks method'.
Value() []byte Value() []byte

View file

@ -329,7 +329,7 @@ func (p *DB) Delete(key []byte) error {
h := p.nodeData[node+nHeight] h := p.nodeData[node+nHeight]
for i, n := range p.prevNode[:h] { for i, n := range p.prevNode[:h] {
m := n + 4 + i m := n + nNext + i
p.nodeData[m] = p.nodeData[p.nodeData[m]+nNext+i] p.nodeData[m] = p.nodeData[p.nodeData[m]+nNext+i]
} }

View file

@ -19,7 +19,7 @@ var (
// Releaser is the interface that wraps the basic Release method. // Releaser is the interface that wraps the basic Release method.
type Releaser interface { type Releaser interface {
// Release releases associated resources. Release should always success // Release releases associated resources. Release should always success
// and can be called multipe times without causing error. // and can be called multiple times without causing error.
Release() Release()
} }

6
vendor/vendor.json vendored
View file

@ -394,10 +394,10 @@
"revisionTime": "2017-07-05T02:17:15Z" "revisionTime": "2017-07-05T02:17:15Z"
}, },
{ {
"checksumSHA1": "yHbyLpI/Meh0DGrmi8x6FrDxxUY=", "checksumSHA1": "l9hsW4atYllqlNjJUwHb6lyLn+I=",
"path": "github.com/syndtr/goleveldb/leveldb", "path": "github.com/syndtr/goleveldb/leveldb",
"revision": "b89cc31ef7977104127d34c1bd31ebd1a9db2199", "revision": "34011bf325bce385408353a30b101fe5e923eb6e",
"revisionTime": "2017-07-25T06:48:36Z" "revisionTime": "2017-12-14T12:08:11Z"
}, },
{ {
"checksumSHA1": "EKIow7XkgNdWvR/982ffIZxKG8Y=", "checksumSHA1": "EKIow7XkgNdWvR/982ffIZxKG8Y=",