mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 14:46:42 +00:00
ethdb: drop concept of cache distribution between dbs
This commit is contained in:
parent
72a0b5572d
commit
59cfdc5302
1 changed files with 60 additions and 76 deletions
|
|
@ -17,7 +17,6 @@
|
||||||
package ethdb
|
package ethdb
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"path/filepath"
|
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
@ -36,20 +35,6 @@ import (
|
||||||
|
|
||||||
var OpenFileLimit = 64
|
var OpenFileLimit = 64
|
||||||
|
|
||||||
// cacheRatio specifies how the total allotted cache is distributed between the
|
|
||||||
// various system databases.
|
|
||||||
var cacheRatio = map[string]float64{
|
|
||||||
"chaindata": 1.0,
|
|
||||||
"lightchaindata": 1.0,
|
|
||||||
}
|
|
||||||
|
|
||||||
// handleRatio specifies how the total allotted file descriptors is distributed
|
|
||||||
// between the various system databases.
|
|
||||||
var handleRatio = map[string]float64{
|
|
||||||
"chaindata": 1.0,
|
|
||||||
"lightchaindata": 1.0,
|
|
||||||
}
|
|
||||||
|
|
||||||
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
|
||||||
|
|
@ -72,16 +57,15 @@ type LDBDatabase struct {
|
||||||
|
|
||||||
// NewLDBDatabase returns a LevelDB wrapped object.
|
// NewLDBDatabase returns a LevelDB wrapped object.
|
||||||
func NewLDBDatabase(file string, cache int, handles int) (*LDBDatabase, error) {
|
func NewLDBDatabase(file string, cache int, handles int) (*LDBDatabase, error) {
|
||||||
// Calculate the cache and file descriptor allowance for this particular database
|
logger := log.New("database", file)
|
||||||
cache = int(float64(cache) * cacheRatio[filepath.Base(file)])
|
|
||||||
|
// Ensure we have some minimal caching and file guarantees
|
||||||
if cache < 16 {
|
if cache < 16 {
|
||||||
cache = 16
|
cache = 16
|
||||||
}
|
}
|
||||||
handles = int(float64(handles) * handleRatio[filepath.Base(file)])
|
|
||||||
if handles < 16 {
|
if handles < 16 {
|
||||||
handles = 16
|
handles = 16
|
||||||
}
|
}
|
||||||
logger := log.New("database", file)
|
|
||||||
logger.Info("Allocated cache and file handles", "cache", cache, "handles", handles)
|
logger.Info("Allocated cache and file handles", "cache", cache, "handles", handles)
|
||||||
|
|
||||||
// Open the db and recover any potential corruptions
|
// Open the db and recover any potential corruptions
|
||||||
|
|
@ -111,103 +95,103 @@ func (db *LDBDatabase) Path() string {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Put puts the given key / value to the queue
|
// Put puts the given key / value to the queue
|
||||||
func (self *LDBDatabase) Put(key []byte, value []byte) error {
|
func (db *LDBDatabase) Put(key []byte, value []byte) error {
|
||||||
// Measure the database put latency, if requested
|
// Measure the database put latency, if requested
|
||||||
if self.putTimer != nil {
|
if db.putTimer != nil {
|
||||||
defer self.putTimer.UpdateSince(time.Now())
|
defer db.putTimer.UpdateSince(time.Now())
|
||||||
}
|
}
|
||||||
// Generate the data to write to disk, update the meter and write
|
// Generate the data to write to disk, update the meter and write
|
||||||
//value = rle.Compress(value)
|
//value = rle.Compress(value)
|
||||||
|
|
||||||
if self.writeMeter != nil {
|
if db.writeMeter != nil {
|
||||||
self.writeMeter.Mark(int64(len(value)))
|
db.writeMeter.Mark(int64(len(value)))
|
||||||
}
|
}
|
||||||
return self.db.Put(key, value, nil)
|
return db.db.Put(key, value, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get returns the given key if it's present.
|
// Get returns the given key if it's present.
|
||||||
func (self *LDBDatabase) Get(key []byte) ([]byte, error) {
|
func (db *LDBDatabase) Get(key []byte) ([]byte, error) {
|
||||||
// Measure the database get latency, if requested
|
// Measure the database get latency, if requested
|
||||||
if self.getTimer != nil {
|
if db.getTimer != nil {
|
||||||
defer self.getTimer.UpdateSince(time.Now())
|
defer db.getTimer.UpdateSince(time.Now())
|
||||||
}
|
}
|
||||||
// Retrieve the key and increment the miss counter if not found
|
// Retrieve the key and increment the miss counter if not found
|
||||||
dat, err := self.db.Get(key, nil)
|
dat, err := db.db.Get(key, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if self.missMeter != nil {
|
if db.missMeter != nil {
|
||||||
self.missMeter.Mark(1)
|
db.missMeter.Mark(1)
|
||||||
}
|
}
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
// Otherwise update the actually retrieved amount of data
|
// Otherwise update the actually retrieved amount of data
|
||||||
if self.readMeter != nil {
|
if db.readMeter != nil {
|
||||||
self.readMeter.Mark(int64(len(dat)))
|
db.readMeter.Mark(int64(len(dat)))
|
||||||
}
|
}
|
||||||
return dat, nil
|
return dat, nil
|
||||||
//return rle.Decompress(dat)
|
//return rle.Decompress(dat)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Delete deletes the key from the queue and database
|
// Delete deletes the key from the queue and database
|
||||||
func (self *LDBDatabase) Delete(key []byte) error {
|
func (db *LDBDatabase) Delete(key []byte) error {
|
||||||
// Measure the database delete latency, if requested
|
// Measure the database delete latency, if requested
|
||||||
if self.delTimer != nil {
|
if db.delTimer != nil {
|
||||||
defer self.delTimer.UpdateSince(time.Now())
|
defer db.delTimer.UpdateSince(time.Now())
|
||||||
}
|
}
|
||||||
// Execute the actual operation
|
// Execute the actual operation
|
||||||
return self.db.Delete(key, nil)
|
return db.db.Delete(key, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *LDBDatabase) NewIterator() iterator.Iterator {
|
func (db *LDBDatabase) NewIterator() iterator.Iterator {
|
||||||
return self.db.NewIterator(nil, nil)
|
return db.db.NewIterator(nil, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *LDBDatabase) Close() {
|
func (db *LDBDatabase) Close() {
|
||||||
// Stop the metrics collection to avoid internal database races
|
// Stop the metrics collection to avoid internal database races
|
||||||
self.quitLock.Lock()
|
db.quitLock.Lock()
|
||||||
defer self.quitLock.Unlock()
|
defer db.quitLock.Unlock()
|
||||||
|
|
||||||
if self.quitChan != nil {
|
if db.quitChan != nil {
|
||||||
errc := make(chan error)
|
errc := make(chan error)
|
||||||
self.quitChan <- errc
|
db.quitChan <- errc
|
||||||
if err := <-errc; err != nil {
|
if err := <-errc; err != nil {
|
||||||
self.log.Error("Metrics collection failed", "err", err)
|
db.log.Error("Metrics collection failed", "err", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
err := self.db.Close()
|
err := db.db.Close()
|
||||||
if err == nil {
|
if err == nil {
|
||||||
self.log.Info("Database closed")
|
db.log.Info("Database closed")
|
||||||
} else {
|
} else {
|
||||||
self.log.Error("Failed to close database", "err", err)
|
db.log.Error("Failed to close database", "err", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *LDBDatabase) LDB() *leveldb.DB {
|
func (db *LDBDatabase) LDB() *leveldb.DB {
|
||||||
return self.db
|
return db.db
|
||||||
}
|
}
|
||||||
|
|
||||||
// Meter configures the database metrics collectors and
|
// Meter configures the database metrics collectors and
|
||||||
func (self *LDBDatabase) Meter(prefix string) {
|
func (db *LDBDatabase) Meter(prefix string) {
|
||||||
// Short circuit metering if the metrics system is disabled
|
// Short circuit metering if the metrics system is disabled
|
||||||
if !metrics.Enabled {
|
if !metrics.Enabled {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Initialize all the metrics collector at the requested prefix
|
// Initialize all the metrics collector at the requested prefix
|
||||||
self.getTimer = metrics.NewTimer(prefix + "user/gets")
|
db.getTimer = metrics.NewTimer(prefix + "user/gets")
|
||||||
self.putTimer = metrics.NewTimer(prefix + "user/puts")
|
db.putTimer = metrics.NewTimer(prefix + "user/puts")
|
||||||
self.delTimer = metrics.NewTimer(prefix + "user/dels")
|
db.delTimer = metrics.NewTimer(prefix + "user/dels")
|
||||||
self.missMeter = metrics.NewMeter(prefix + "user/misses")
|
db.missMeter = metrics.NewMeter(prefix + "user/misses")
|
||||||
self.readMeter = metrics.NewMeter(prefix + "user/reads")
|
db.readMeter = metrics.NewMeter(prefix + "user/reads")
|
||||||
self.writeMeter = metrics.NewMeter(prefix + "user/writes")
|
db.writeMeter = metrics.NewMeter(prefix + "user/writes")
|
||||||
self.compTimeMeter = metrics.NewMeter(prefix + "compact/time")
|
db.compTimeMeter = metrics.NewMeter(prefix + "compact/time")
|
||||||
self.compReadMeter = metrics.NewMeter(prefix + "compact/input")
|
db.compReadMeter = metrics.NewMeter(prefix + "compact/input")
|
||||||
self.compWriteMeter = metrics.NewMeter(prefix + "compact/output")
|
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
|
||||||
self.quitLock.Lock()
|
db.quitLock.Lock()
|
||||||
self.quitChan = make(chan chan error)
|
db.quitChan = make(chan chan error)
|
||||||
self.quitLock.Unlock()
|
db.quitLock.Unlock()
|
||||||
|
|
||||||
go self.meter(3 * time.Second)
|
go db.meter(3 * time.Second)
|
||||||
}
|
}
|
||||||
|
|
||||||
// meter periodically retrieves internal leveldb counters and reports them to
|
// meter periodically retrieves internal leveldb counters and reports them to
|
||||||
|
|
@ -221,7 +205,7 @@ func (self *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 (self *LDBDatabase) meter(refresh time.Duration) {
|
func (db *LDBDatabase) meter(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++ {
|
||||||
|
|
@ -230,9 +214,9 @@ func (self *LDBDatabase) meter(refresh time.Duration) {
|
||||||
// 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
|
// Retrieve the database stats
|
||||||
stats, err := self.db.GetProperty("leveldb.stats")
|
stats, err := db.db.GetProperty("leveldb.stats")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
self.log.Error("Failed to read database stats", "err", err)
|
db.log.Error("Failed to read database stats", "err", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Find the compaction table, skip the header
|
// Find the compaction table, skip the header
|
||||||
|
|
@ -241,7 +225,7 @@ func (self *LDBDatabase) meter(refresh time.Duration) {
|
||||||
lines = lines[1:]
|
lines = lines[1:]
|
||||||
}
|
}
|
||||||
if len(lines) <= 3 {
|
if len(lines) <= 3 {
|
||||||
self.log.Error("Compaction table not found")
|
db.log.Error("Compaction table not found")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
lines = lines[3:]
|
lines = lines[3:]
|
||||||
|
|
@ -258,25 +242,25 @@ func (self *LDBDatabase) meter(refresh time.Duration) {
|
||||||
for idx, counter := range parts[3:] {
|
for idx, counter := range parts[3:] {
|
||||||
value, err := strconv.ParseFloat(strings.TrimSpace(counter), 64)
|
value, err := strconv.ParseFloat(strings.TrimSpace(counter), 64)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
self.log.Error("Compaction entry parsing failed", "err", err)
|
db.log.Error("Compaction entry parsing failed", "err", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
counters[i%2][idx] += value
|
counters[i%2][idx] += value
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Update all the requested meters
|
// Update all the requested meters
|
||||||
if self.compTimeMeter != nil {
|
if db.compTimeMeter != nil {
|
||||||
self.compTimeMeter.Mark(int64((counters[i%2][0] - counters[(i-1)%2][0]) * 1000 * 1000 * 1000))
|
db.compTimeMeter.Mark(int64((counters[i%2][0] - counters[(i-1)%2][0]) * 1000 * 1000 * 1000))
|
||||||
}
|
}
|
||||||
if self.compReadMeter != nil {
|
if db.compReadMeter != nil {
|
||||||
self.compReadMeter.Mark(int64((counters[i%2][1] - counters[(i-1)%2][1]) * 1024 * 1024))
|
db.compReadMeter.Mark(int64((counters[i%2][1] - counters[(i-1)%2][1]) * 1024 * 1024))
|
||||||
}
|
}
|
||||||
if self.compWriteMeter != nil {
|
if db.compWriteMeter != nil {
|
||||||
self.compWriteMeter.Mark(int64((counters[i%2][2] - counters[(i-1)%2][2]) * 1024 * 1024))
|
db.compWriteMeter.Mark(int64((counters[i%2][2] - counters[(i-1)%2][2]) * 1024 * 1024))
|
||||||
}
|
}
|
||||||
// Sleep a bit, then repeat the stats collection
|
// Sleep a bit, then repeat the stats collection
|
||||||
select {
|
select {
|
||||||
case errc := <-self.quitChan:
|
case errc := <-db.quitChan:
|
||||||
// Quit requesting, stop hammering the database
|
// Quit requesting, stop hammering the database
|
||||||
errc <- nil
|
errc <- nil
|
||||||
return
|
return
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue