diff --git a/cmd/geth/chaincmd.go b/cmd/geth/chaincmd.go index 85d0c3acaa..7adfea5285 100644 --- a/cmd/geth/chaincmd.go +++ b/cmd/geth/chaincmd.go @@ -225,6 +225,11 @@ func importChain(ctx *cli.Context) error { utils.Fatalf("Failed to read database stats: %v", err) } fmt.Println(stats) + iostats, err := db.LDB().GetProperty("leveldb.iostats") + if err != nil { + utils.Fatalf("Failed to read database iostats: %v", err) + } + fmt.Println(iostats) fmt.Printf("Trie cache misses: %d\n", trie.CacheMisses()) fmt.Printf("Trie cache unloads: %d\n\n", trie.CacheUnloads()) @@ -254,6 +259,11 @@ func importChain(ctx *cli.Context) error { utils.Fatalf("Failed to read database stats: %v", err) } fmt.Println(stats) + iostats, err = db.LDB().GetProperty("leveldb.iostats") + if err != nil { + utils.Fatalf("Failed to read database iostats: %v", err) + } + fmt.Println(iostats) return nil } diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index 2ca795187f..0193794778 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -281,8 +281,8 @@ func (db *Dashboard) collectData() { prevNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count() prevProcessCPUTime = getProcessCPUTime() prevSystemCPUUsage = systemCPUUsage - prevDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/input").(metrics.Meter).Count() - prevDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/output").(metrics.Meter).Count() + prevDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/read").(metrics.Meter).Count() + prevDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/write").(metrics.Meter).Count() frequency = float64(db.config.Refresh / time.Second) numCPU = float64(runtime.NumCPU()) @@ -300,8 +300,8 @@ func (db *Dashboard) collectData() { curNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count() curProcessCPUTime = getProcessCPUTime() curSystemCPUUsage = systemCPUUsage - curDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/input").(metrics.Meter).Count() - curDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/output").(metrics.Meter).Count() + curDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/read").(metrics.Meter).Count() + curDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/write").(metrics.Meter).Count() deltaNetworkIngress = float64(curNetworkIngress - prevNetworkIngress) deltaNetworkEgress = float64(curNetworkEgress - prevNetworkEgress) diff --git a/ethdb/database.go b/ethdb/database.go index 57d38f7f5f..32f26258ad 100644 --- a/ethdb/database.go +++ b/ethdb/database.go @@ -22,6 +22,7 @@ import ( "sync" "time" + "fmt" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/metrics" "github.com/syndtr/goleveldb/leveldb" @@ -29,6 +30,7 @@ import ( "github.com/syndtr/goleveldb/leveldb/filter" "github.com/syndtr/goleveldb/leveldb/iterator" "github.com/syndtr/goleveldb/leveldb/opt" + "regexp" ) var OpenFileLimit = 64 @@ -46,6 +48,8 @@ type LDBDatabase struct { compTimeMeter metrics.Meter // Meter for measuring the total time spent in database compaction compReadMeter metrics.Meter // Meter for measuring the data read during compaction compWriteMeter metrics.Meter // Meter for measuring the data written during compaction + diskReadMeter metrics.Meter // Meter for measuring the effective amount of data read + diskWriteMeter metrics.Meter // Meter for measuring the effective amount of data written quitLock sync.Mutex // Mutex protecting the quit channel access quitChan chan chan error // Quit channel to stop the metrics collection before closing the database @@ -187,16 +191,67 @@ func (db *LDBDatabase) Meter(prefix string) { db.compTimeMeter = metrics.NewRegisteredMeter(prefix+"compact/time", nil) db.compReadMeter = metrics.NewRegisteredMeter(prefix+"compact/input", nil) db.compWriteMeter = metrics.NewRegisteredMeter(prefix+"compact/output", nil) + db.diskReadMeter = metrics.NewRegisteredMeter(prefix+"disk/read", nil) + db.diskWriteMeter = metrics.NewRegisteredMeter(prefix+"disk/write", nil) // Create a quit channel for the periodic collector and run it db.quitLock.Lock() db.quitChan = make(chan chan error) db.quitLock.Unlock() - go db.meter(3 * time.Second) + go db.meterDiskIO(3 * time.Second) + go db.meterCompaction(3 * time.Second) } -// meter periodically retrieves internal leveldb counters and reports them to +// meterDiskIO periodically retrieves internal leveldb counters and reports them to +// the metrics subsystem. +// +// This is how the iostats look like (currently): +// Read(MB): 3895.04860 Write(MB): 3654.64712 +func (db *LDBDatabase) meterDiskIO(refresh time.Duration) { + var prev, curr [2]float64 + spaceTruncater := regexp.MustCompile(`[ ]{2,}`) + + for { + ioStats, err := db.db.GetProperty("leveldb.iostats") + if err != nil { + db.log.Error("Failed to read database iostats", "err", err) + return + } + + parts := strings.Split(spaceTruncater.ReplaceAllString(ioStats, " "), " ") + if curr[0], err = strconv.ParseFloat(parts[1], 64); err != nil { + db.log.Error("Read entry parsing failed", "err", err) + return + } + if curr[1], err = strconv.ParseFloat(parts[3], 64); err != nil { + db.log.Error("Write entry parsing failed", "err", err) + return + } + if db.diskReadMeter != nil { + db.diskReadMeter.Mark(int64((curr[0] - prev[0]) * 1024 * 1024)) + } + if db.diskWriteMeter != nil { + db.diskWriteMeter.Mark(int64((curr[1] - prev[1]) * 1024 * 1024)) + } + fmt.Printf("Read: %9.5fMB Write: %9.5fMB / %v\n", curr[0]-prev[0], curr[1]-prev[1], refresh) + prev[0] = curr[0] + prev[1] = curr[1] + + // Sleep a bit, then repeat the iostats collection + select { + case errc := <-db.quitChan: + // Quit requesting, stop hammering the database + errc <- nil + return + + case <-time.After(refresh): + // Timeout, gather a new set of iostats + } + } +} + +// meterCompaction periodically retrieves internal leveldb counters and reports them to // the metrics subsystem. // // This is how a stats table look like (currently): @@ -207,7 +262,7 @@ func (db *LDBDatabase) Meter(prefix string) { // 1 | 85 | 109.27913 | 28.09293 | 213.92493 | 214.26294 // 2 | 523 | 1000.37159 | 7.26059 | 66.86342 | 66.77884 // 3 | 570 | 1113.18458 | 0.00000 | 0.00000 | 0.00000 -func (db *LDBDatabase) meter(refresh time.Duration) { +func (db *LDBDatabase) meterCompaction(refresh time.Duration) { // Create the counters to store current and previous values counters := make([][]float64, 2) for i := 0; i < 2; i++ { diff --git a/vendor/github.com/syndtr/goleveldb/leveldb/db.go b/vendor/github.com/syndtr/goleveldb/leveldb/db.go index ea5595eb3a..ff18235141 100644 --- a/vendor/github.com/syndtr/goleveldb/leveldb/db.go +++ b/vendor/github.com/syndtr/goleveldb/leveldb/db.go @@ -168,7 +168,7 @@ func openDB(s *session) (*DB, error) { // The returned DB instance is safe for concurrent use. // The DB must be closed after use, by calling Close method. func Open(stor storage.Storage, o *opt.Options) (db *DB, err error) { - s, err := newSession(stor, o) + s, err := newSession(storage.IOCounterWrapper(stor), o) if err != nil { return } @@ -906,6 +906,8 @@ func (db *DB) GetSnapshot() (*Snapshot, error) { // Returns the number of files at level 'n'. // leveldb.stats // Returns statistics of the underlying DB. +// leveldb.iostats +// Returns statistics of effective disk read and write. // leveldb.writedelay // Returns cumulative write delay caused by compaction. // leveldb.sstables @@ -959,6 +961,12 @@ func (db *DB) GetProperty(name string) (value string, err error) { level, len(tables), float64(tables.size())/1048576.0, duration.Seconds(), float64(read)/1048576.0, float64(write)/1048576.0) } + case p == "iostats": + var r, w float64 + if s, ok := db.s.stor.(storage.IOCounter); ok { + r, w = float64(s.Reads())/1048576.0, float64(s.Writes())/1048576.0 + } + value = fmt.Sprintf("Read(MB): %13.5f Write(MB): %13.5f", r, w) case p == "writedelay": writeDelayN, writeDelay := atomic.LoadInt32(&db.cWriteDelayN), time.Duration(atomic.LoadInt64(&db.cWriteDelay)) value = fmt.Sprintf("DelayN:%d Delay:%s", writeDelayN, writeDelay) diff --git a/vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go b/vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go index c16bce6b66..301a484ecf 100644 --- a/vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go +++ b/vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go @@ -11,6 +11,7 @@ import ( "errors" "fmt" "io" + "sync/atomic" ) // FileType represent a file type. @@ -177,3 +178,78 @@ type Storage interface { // called after the storage has been closed. Close() error } + +// IOCounter collects read and write statistics. +type IOCounter interface { + // Reads returns the cumulative number of read bytes of the underlying storage. + Reads() uint64 + // Writes returns the cumulative number of written bytes of the underlying storage. + Writes() uint64 +} + +type ioCounter struct { + Storage + read uint64 + write uint64 +} + +func (c *ioCounter) Open(fd FileDesc) (Reader, error) { + r, err := c.Storage.Open(fd) + return &meteredReader{r, c}, err +} + +func (c *ioCounter) Create(fd FileDesc) (Writer, error) { + w, err := c.Storage.Create(fd) + return &meteredWriter{w, c}, err +} + +func (c *ioCounter) Reads() uint64 { + return atomic.LoadUint64(&c.read) +} + +func (c *ioCounter) Writes() uint64 { + return atomic.LoadUint64(&c.write) +} + +// AddRead increases the number of read bytes by n. +func (c *ioCounter) AddRead(n uint64) uint64 { + return atomic.AddUint64(&c.read, n) +} + +// AddWrite increases the number of written bytes by n. +func (c *ioCounter) AddWrite(n uint64) uint64 { + return atomic.AddUint64(&c.write, n) +} + +// IOCounterWrapper returns the given storage wrapped by ioCounter. +func IOCounterWrapper(s Storage) Storage { + return &ioCounter{s, 0, 0} +} + +type meteredReader struct { + Reader + c *ioCounter +} + +func (r *meteredReader) Read(p []byte) (n int, err error) { + n, err = r.Reader.Read(p) + r.c.AddRead(uint64(n)) + return n, err +} + +func (r *meteredReader) ReadAt(p []byte, off int64) (n int, err error) { + n, err = r.Reader.ReadAt(p, off) + r.c.AddRead(uint64(n)) + return n, err +} + +type meteredWriter struct { + Writer + c *ioCounter +} + +func (w *meteredWriter) Write(p []byte) (n int, err error) { + n, err = w.Writer.Write(p) + w.c.AddWrite(uint64(n)) + return n, err +}