mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 18:02:24 +00:00
Merge 6a614171a8 into 3ec1b9a92d
This commit is contained in:
commit
c7818400e8
5 changed files with 157 additions and 8 deletions
|
|
@ -225,6 +225,11 @@ func importChain(ctx *cli.Context) error {
|
||||||
utils.Fatalf("Failed to read database stats: %v", err)
|
utils.Fatalf("Failed to read database stats: %v", err)
|
||||||
}
|
}
|
||||||
fmt.Println(stats)
|
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 misses: %d\n", trie.CacheMisses())
|
||||||
fmt.Printf("Trie cache unloads: %d\n\n", trie.CacheUnloads())
|
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)
|
utils.Fatalf("Failed to read database stats: %v", err)
|
||||||
}
|
}
|
||||||
fmt.Println(stats)
|
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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -281,8 +281,8 @@ func (db *Dashboard) collectData() {
|
||||||
prevNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count()
|
prevNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count()
|
||||||
prevProcessCPUTime = getProcessCPUTime()
|
prevProcessCPUTime = getProcessCPUTime()
|
||||||
prevSystemCPUUsage = systemCPUUsage
|
prevSystemCPUUsage = systemCPUUsage
|
||||||
prevDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/input").(metrics.Meter).Count()
|
prevDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/read").(metrics.Meter).Count()
|
||||||
prevDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/output").(metrics.Meter).Count()
|
prevDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/write").(metrics.Meter).Count()
|
||||||
|
|
||||||
frequency = float64(db.config.Refresh / time.Second)
|
frequency = float64(db.config.Refresh / time.Second)
|
||||||
numCPU = float64(runtime.NumCPU())
|
numCPU = float64(runtime.NumCPU())
|
||||||
|
|
@ -300,8 +300,8 @@ func (db *Dashboard) collectData() {
|
||||||
curNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count()
|
curNetworkEgress = metrics.DefaultRegistry.Get("p2p/OutboundTraffic").(metrics.Meter).Count()
|
||||||
curProcessCPUTime = getProcessCPUTime()
|
curProcessCPUTime = getProcessCPUTime()
|
||||||
curSystemCPUUsage = systemCPUUsage
|
curSystemCPUUsage = systemCPUUsage
|
||||||
curDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/input").(metrics.Meter).Count()
|
curDiskRead = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/read").(metrics.Meter).Count()
|
||||||
curDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/compact/output").(metrics.Meter).Count()
|
curDiskWrite = metrics.DefaultRegistry.Get("eth/db/chaindata/disk/write").(metrics.Meter).Count()
|
||||||
|
|
||||||
deltaNetworkIngress = float64(curNetworkIngress - prevNetworkIngress)
|
deltaNetworkIngress = float64(curNetworkIngress - prevNetworkIngress)
|
||||||
deltaNetworkEgress = float64(curNetworkEgress - prevNetworkEgress)
|
deltaNetworkEgress = float64(curNetworkEgress - prevNetworkEgress)
|
||||||
|
|
|
||||||
|
|
@ -22,6 +22,7 @@ import (
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"fmt"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/metrics"
|
"github.com/ethereum/go-ethereum/metrics"
|
||||||
"github.com/syndtr/goleveldb/leveldb"
|
"github.com/syndtr/goleveldb/leveldb"
|
||||||
|
|
@ -29,6 +30,7 @@ import (
|
||||||
"github.com/syndtr/goleveldb/leveldb/filter"
|
"github.com/syndtr/goleveldb/leveldb/filter"
|
||||||
"github.com/syndtr/goleveldb/leveldb/iterator"
|
"github.com/syndtr/goleveldb/leveldb/iterator"
|
||||||
"github.com/syndtr/goleveldb/leveldb/opt"
|
"github.com/syndtr/goleveldb/leveldb/opt"
|
||||||
|
"regexp"
|
||||||
)
|
)
|
||||||
|
|
||||||
var OpenFileLimit = 64
|
var OpenFileLimit = 64
|
||||||
|
|
@ -46,6 +48,8 @@ type LDBDatabase struct {
|
||||||
compTimeMeter metrics.Meter // Meter for measuring the total time spent in database compaction
|
compTimeMeter metrics.Meter // Meter for measuring the total time spent in database compaction
|
||||||
compReadMeter metrics.Meter // Meter for measuring the data read during compaction
|
compReadMeter metrics.Meter // Meter for measuring the data read during compaction
|
||||||
compWriteMeter metrics.Meter // Meter for measuring the data written 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
|
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
|
||||||
|
|
@ -187,16 +191,67 @@ func (db *LDBDatabase) Meter(prefix string) {
|
||||||
db.compTimeMeter = metrics.NewRegisteredMeter(prefix+"compact/time", nil)
|
db.compTimeMeter = metrics.NewRegisteredMeter(prefix+"compact/time", nil)
|
||||||
db.compReadMeter = metrics.NewRegisteredMeter(prefix+"compact/input", nil)
|
db.compReadMeter = metrics.NewRegisteredMeter(prefix+"compact/input", nil)
|
||||||
db.compWriteMeter = metrics.NewRegisteredMeter(prefix+"compact/output", 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
|
// Create a quit channel for the periodic collector and run it
|
||||||
db.quitLock.Lock()
|
db.quitLock.Lock()
|
||||||
db.quitChan = make(chan chan error)
|
db.quitChan = make(chan chan error)
|
||||||
db.quitLock.Unlock()
|
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.
|
// the metrics subsystem.
|
||||||
//
|
//
|
||||||
// This is how a stats table look like (currently):
|
// 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
|
// 1 | 85 | 109.27913 | 28.09293 | 213.92493 | 214.26294
|
||||||
// 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) meterCompaction(refresh time.Duration) {
|
||||||
// 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++ {
|
||||||
|
|
|
||||||
10
vendor/github.com/syndtr/goleveldb/leveldb/db.go
generated
vendored
10
vendor/github.com/syndtr/goleveldb/leveldb/db.go
generated
vendored
|
|
@ -168,7 +168,7 @@ func openDB(s *session) (*DB, error) {
|
||||||
// The returned DB instance is safe for concurrent use.
|
// The returned DB instance is safe for concurrent use.
|
||||||
// The DB must be closed after use, by calling Close method.
|
// The DB must be closed after use, by calling Close method.
|
||||||
func Open(stor storage.Storage, o *opt.Options) (db *DB, err error) {
|
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 {
|
if err != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -906,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.iostats
|
||||||
|
// Returns statistics of effective disk read and write.
|
||||||
// leveldb.writedelay
|
// leveldb.writedelay
|
||||||
// Returns cumulative write delay caused by compaction.
|
// Returns cumulative write delay caused by compaction.
|
||||||
// leveldb.sstables
|
// 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(),
|
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 == "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":
|
case p == "writedelay":
|
||||||
writeDelayN, writeDelay := atomic.LoadInt32(&db.cWriteDelayN), time.Duration(atomic.LoadInt64(&db.cWriteDelay))
|
writeDelayN, writeDelay := atomic.LoadInt32(&db.cWriteDelayN), time.Duration(atomic.LoadInt64(&db.cWriteDelay))
|
||||||
value = fmt.Sprintf("DelayN:%d Delay:%s", writeDelayN, writeDelay)
|
value = fmt.Sprintf("DelayN:%d Delay:%s", writeDelayN, writeDelay)
|
||||||
|
|
|
||||||
76
vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go
generated
vendored
76
vendor/github.com/syndtr/goleveldb/leveldb/storage/storage.go
generated
vendored
|
|
@ -11,6 +11,7 @@ import (
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"sync/atomic"
|
||||||
)
|
)
|
||||||
|
|
||||||
// FileType represent a file type.
|
// FileType represent a file type.
|
||||||
|
|
@ -177,3 +178,78 @@ type Storage interface {
|
||||||
// called after the storage has been closed.
|
// called after the storage has been closed.
|
||||||
Close() error
|
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
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue