From 4f0cc2d7e23c0973b2bf2ddba917b3525936003c Mon Sep 17 00:00:00 2001 From: rjl493456442 Date: Thu, 18 Jul 2019 15:19:31 +0800 Subject: [PATCH] core: init tx lookup in background --- core/blockchain.go | 8 +- core/rawdb/accessors_chain.go | 26 +++++ core/rawdb/{freezer_reinit.go => initer.go} | 101 +++++++++++++++----- core/rawdb/schema.go | 3 + 4 files changed, 113 insertions(+), 25 deletions(-) rename core/rawdb/{freezer_reinit.go => initer.go} (51%) diff --git a/core/blockchain.go b/core/blockchain.go index 833de3bc7e..8d8476df38 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -230,7 +230,13 @@ func NewBlockChain(db ethdb.Database, cacheConfig *CacheConfig, chainConfig *par } // Initialize the chain with ancient data if it isn't empty. if bc.empty() { - rawdb.InitDatabaseFromFreezer(bc.db) + rawdb.InitBlockIndexFromFreezer(bc.db) + rawdb.WriteAncientTxLookupProgress(bc.db, 0) // Explicitly mark the missing of txlookup. + } + // Re-initialise all ancient txlookup indexes in the background. + if number := rawdb.ReadAncientTxLookupProgress(bc.db); number != nil { + // Genesis block doesn't have transaction, just ignore it. + go rawdb.InitTxsLookupFromFreezer(bc.db, *number+1) } if err := bc.loadLastState(); err != nil { return nil, err diff --git a/core/rawdb/accessors_chain.go b/core/rawdb/accessors_chain.go index fab7ca56c3..50df9adc79 100644 --- a/core/rawdb/accessors_chain.go +++ b/core/rawdb/accessors_chain.go @@ -171,6 +171,32 @@ func WriteFastTrieProgress(db ethdb.KeyValueWriter, count uint64) { } } +// ReadAncientTxLookupProgress retrieves the number of ancient blocks which +// txlookup has been inserted to allow reporting correct numbers across restarts. +func ReadAncientTxLookupProgress(db ethdb.KeyValueReader) *uint64 { + data, _ := db.Get(ancientTxLookupProgressKey) + if len(data) != 8 { + return nil + } + number := binary.BigEndian.Uint64(data) + return &number +} + +// WriteAncientTxLookupProgress stores the ancient txlookup process counter to support +// retrieving it across restarts. +func WriteAncientTxLookupProgress(db ethdb.KeyValueWriter, number uint64) { + if err := db.Put(ancientTxLookupProgressKey, encodeBlockNumber(number)); err != nil { + log.Crit("Failed to store head number of txlookup", "err", err) + } +} + +// DeleteAncientTxLookupProgress deletes the ancient txlookup progress. +func DeleteAncientTxLookupProgress(db ethdb.KeyValueWriter) { + if err := db.Delete(ancientTxLookupProgressKey); err != nil { + log.Crit("Failed to delete ancient txlookup progress entry", "err", err) + } +} + // ReadHeaderRLP retrieves a block header in its raw RLP database encoding. func ReadHeaderRLP(db ethdb.Reader, hash common.Hash, number uint64) rlp.RawValue { data, _ := db.Ancient(freezerHeaderTable, number) diff --git a/core/rawdb/freezer_reinit.go b/core/rawdb/initer.go similarity index 51% rename from core/rawdb/freezer_reinit.go rename to core/rawdb/initer.go index ea4dd33d1d..d5116d558d 100644 --- a/core/rawdb/freezer_reinit.go +++ b/core/rawdb/initer.go @@ -29,21 +29,34 @@ import ( "github.com/ethereum/go-ethereum/log" ) -// InitDatabaseFromFreezer reinitializes an empty database from a previous batch -// of frozen ancient blocks. The method iterates over all the frozen blocks and -// injects into the database the block hash->number mappings and the transaction -// lookup entries. -func InitDatabaseFromFreezer(db ethdb.Database) error { +type ( + initPrepare func(*types.Block) // The callback for customized prepare operation. + initAction func(ethdb.Batch, *types.Block) // The callback for customized initialisation action. +) + +// iterateAncient iterates the specified range blocks from ancient database +// and then apply initialisation action. +func iterateAncient(db ethdb.Database, from uint64, typ string, prepare initPrepare, action initAction) error { + // Short circuit if the init action is nil. + if action == nil { + return nil + } // If we can't access the freezer or it's empty, abort frozen, err := db.Ancients() if err != nil || frozen == 0 { return err } - // Blocks previously frozen, iterate over- and hash them concurrently + // Spawn multi-routines, iterate over the specified blocks and invoke prepare + // callback concurrently. var ( - number = ^uint64(0) // -1 + number uint64 results = make(chan *types.Block, 4*runtime.NumCPU()) ) + if from == 0 { + number = ^uint64(0) // -1 + } else { + number = from - 1 + } abort := make(chan struct{}) defer close(abort) @@ -56,14 +69,10 @@ func InitDatabaseFromFreezer(db ethdb.Database) error { return } // Retrieve the block from the freezer (no need for the hash, we pull by - // number from the freezer). If successful, pre-cache the block hash and - // the individual transaction hashes for storing into the database. + // number from the freezer). block := ReadBlock(db, common.Hash{}, n) - if block != nil { - block.Hash() - for _, tx := range block.Transactions() { - tx.Hash() - } + if prepare != nil && block != nil { + prepare(block) } // Feed the block to the aggregator, or abort on interrupt select { @@ -74,20 +83,20 @@ func InitDatabaseFromFreezer(db ethdb.Database) error { } }() } - // Reassemble the blocks into a contiguous stream and push them out to disk + // Reassemble the blocks into a contiguous stream and apply the action callback. var ( queue = prque.New(nil) - next = int64(0) + next = int64(from) batch = db.NewBatch() start = time.Now() logged time.Time ) - for i := uint64(0); i < frozen; i++ { + for i := from; i < frozen; i++ { // Retrieve the next result and bail if it's nil block := <-results if block == nil { - return errors.New("broken ancient database") + return errors.New("broken database") } // Push the block into the import queue and process contiguous ranges queue.Push(block, -int64(block.NumberU64())) @@ -100,9 +109,8 @@ func InitDatabaseFromFreezer(db ethdb.Database) error { block = queue.PopItem().(*types.Block) next++ - // Inject hash<->number mapping and txlookup indexes - WriteHeaderNumber(batch, block.Hash(), block.NumberU64()) - WriteTxLookupEntries(batch, block) + // Invoke action to inject specified data into key-value database. + action(batch, block) // If enough data was accumulated in memory or we're at the last block, dump to disk if batch.ValueSize() > ethdb.IdealBatchSize || uint64(next) == frozen { @@ -113,15 +121,60 @@ func InitDatabaseFromFreezer(db ethdb.Database) error { } // If we've spent too much time already, notify the user of what we're doing if time.Since(logged) > 8*time.Second { - log.Info("Initializing chain from ancient data", "number", block.Number(), "hash", block.Hash(), "total", frozen-1, "elapsed", common.PrettyDuration(time.Since(start))) + log.Info("Initializing chain from ancient data", "type", typ, "number", block.Number(), "hash", block.Hash(), "total", uint64(next)-from, "elapsed", common.PrettyDuration(time.Since(start))) logged = time.Now() } } } + log.Info("Initialized chain from ancient data", "type", typ, "number", frozen-from, "elapsed", common.PrettyDuration(time.Since(start))) + return nil +} + +// InitBlockIndexFromFreezer reinitializes an empty database from a previous batch +// of frozen ancient blocks. The method iterates over all the frozen blocks and +// injects into the database the block hash->number mappings and the transaction +// lookup entries. +func InitBlockIndexFromFreezer(db ethdb.Database) error { + // If we can't access the freezer or it's empty, abort + frozen, err := db.Ancients() + if err != nil || frozen == 0 { + return err + } + // hashBlock calculates block hash in advance using the multi-routine's concurrent + // computing power. + hashBlock := func(block *types.Block) { block.Hash() } + + // writeIndex injects hash <-> number mapping into the database. + writeIndex := func(batch ethdb.Batch, block *types.Block) { WriteHeaderNumber(batch, block.Hash(), block.NumberU64()) } + + if err := iterateAncient(db, 0, "blocks", hashBlock, writeIndex); err != nil { + return err + } hash := ReadCanonicalHash(db, frozen-1) WriteHeadHeaderHash(db, hash) WriteHeadFastBlockHash(db, hash) - - log.Info("Initialized chain from ancient data", "number", frozen-1, "hash", hash, "elapsed", common.PrettyDuration(time.Since(start))) + return nil +} + +// InitTxsLookupFromFreezer initializes txlookup indexes in the database. +func InitTxsLookupFromFreezer(db ethdb.Database, from uint64) error { + // hashTxs calculates transaction hash in advance using the multi-routine's + // concurrent computing power. + hashTxs := func(block *types.Block) { + for _, tx := range block.Transactions() { + tx.Hash() + } + } + // writeIndex injects txlookup indexes into the database. + writeIndex := func(batch ethdb.Batch, block *types.Block) { + WriteTxLookupEntries(batch, block) + if block.NumberU64()%10000 == 0 { + WriteAncientTxLookupProgress(batch, block.NumberU64()) + } + } + if err := iterateAncient(db, from, "txlookup", hashTxs, writeIndex); err != nil { + return err + } + DeleteAncientTxLookupProgress(db) // Mark all txlookup indexes of ancient blocks have been inserted. return nil } diff --git a/core/rawdb/schema.go b/core/rawdb/schema.go index a44a2c99f9..34f47c6422 100644 --- a/core/rawdb/schema.go +++ b/core/rawdb/schema.go @@ -41,6 +41,9 @@ var ( // fastTrieProgressKey tracks the number of trie entries imported during fast sync. fastTrieProgressKey = []byte("TrieSync") + // ancientTxLookupProgressKey tracks the progress of ancient txs lookup insertion. + ancientTxLookupProgressKey = []byte("AncientTxsLookup") + // Data item prefixes (use single byte to avoid mixing data types, avoid `i`, used for indexes). headerPrefix = []byte("h") // headerPrefix + num (uint64 big endian) + hash -> header headerTDSuffix = []byte("t") // headerPrefix + num (uint64 big endian) + hash + headerTDSuffix -> td