mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-15 16:33:47 +00:00
Add ability to calculate the longest execution path in a block
This commit is contained in:
parent
ad658b6b99
commit
6f0d16fbeb
4 changed files with 110 additions and 47 deletions
|
|
@ -2,8 +2,8 @@ package blockstm
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"sort"
|
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/heimdalr/dag"
|
"github.com/heimdalr/dag"
|
||||||
|
|
||||||
|
|
@ -62,8 +62,6 @@ func BuildDAG(deps TxnInputOutput) (d DAG) {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warn("Failed to add edge", "from", txFromId, "to", txToId, "err", err)
|
log.Warn("Failed to add edge", "from", txFromId, "to", txToId, "err", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
break // once we add a 'backward' dep we can't execute before that transaction so no need to proceed
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -71,23 +69,67 @@ func BuildDAG(deps TxnInputOutput) (d DAG) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d DAG) Report(out func(string)) {
|
// Find the longest execution path in the DAG
|
||||||
roots := make([]int, 0)
|
func (d DAG) LongestPath(stats map[int]ExecutionStat) ([]int, uint64) {
|
||||||
rootIds := make([]string, 0)
|
prev := make(map[int]int, len(d.GetVertices()))
|
||||||
rootIdMap := make(map[int]string, len(d.GetRoots()))
|
|
||||||
|
|
||||||
for k, i := range d.GetRoots() {
|
for i := 0; i < len(d.GetVertices()); i++ {
|
||||||
roots = append(roots, i.(int))
|
prev[i] = -1
|
||||||
rootIdMap[i.(int)] = k
|
|
||||||
}
|
}
|
||||||
|
|
||||||
sort.Ints(roots)
|
pathWeights := make(map[int]uint64, len(d.GetVertices()))
|
||||||
|
|
||||||
for _, i := range roots {
|
maxPath := 0
|
||||||
rootIds = append(rootIds, rootIdMap[i])
|
maxPathWeight := uint64(0)
|
||||||
|
|
||||||
|
idxToId := make(map[int]string, len(d.GetVertices()))
|
||||||
|
|
||||||
|
for k, i := range d.GetVertices() {
|
||||||
|
idxToId[i.(int)] = k
|
||||||
}
|
}
|
||||||
|
|
||||||
fmt.Println(roots)
|
for i := 0; i < len(idxToId); i++ {
|
||||||
|
parents, _ := d.GetParents(idxToId[i])
|
||||||
|
|
||||||
|
if len(parents) > 0 {
|
||||||
|
for _, p := range parents {
|
||||||
|
weight := pathWeights[p.(int)] + stats[i].End - stats[i].Start
|
||||||
|
if weight > pathWeights[i] {
|
||||||
|
pathWeights[i] = weight
|
||||||
|
prev[i] = p.(int)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
pathWeights[i] = stats[i].End - stats[i].Start
|
||||||
|
}
|
||||||
|
|
||||||
|
if pathWeights[i] > maxPathWeight {
|
||||||
|
maxPath = i
|
||||||
|
maxPathWeight = pathWeights[i]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
path := make([]int, 0)
|
||||||
|
for i := maxPath; i != -1; i = prev[i] {
|
||||||
|
path = append(path, i)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reverse the path so the transactions are in the ascending order
|
||||||
|
for i, j := 0, len(path)-1; i < j; i, j = i+1, j-1 {
|
||||||
|
path[i], path[j] = path[j], path[i]
|
||||||
|
}
|
||||||
|
|
||||||
|
return path, maxPathWeight
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d DAG) Report(stats map[int]ExecutionStat, out func(string)) {
|
||||||
|
longestPath, weight := d.LongestPath(stats)
|
||||||
|
|
||||||
|
serialWeight := uint64(0)
|
||||||
|
|
||||||
|
for i := 0; i < len(d.GetVertices()); i++ {
|
||||||
|
serialWeight += stats[i].End - stats[i].Start
|
||||||
|
}
|
||||||
|
|
||||||
makeStrs := func(ints []int) (ret []string) {
|
makeStrs := func(ints []int) (ret []string) {
|
||||||
for _, v := range ints {
|
for _, v := range ints {
|
||||||
|
|
@ -97,29 +139,9 @@ func (d DAG) Report(out func(string)) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
maxDesc := 0
|
out("Longest execution path:")
|
||||||
maxDeps := 0
|
out(fmt.Sprintf("(%v) %v", len(longestPath), strings.Join(makeStrs(longestPath), "->")))
|
||||||
totalDeps := 0
|
|
||||||
|
|
||||||
for k, v := range roots {
|
out(fmt.Sprintf("Longest path ideal execution time: %v of %v (serial total), %v%%", time.Duration(weight),
|
||||||
ids := []int{v}
|
time.Duration(serialWeight), fmt.Sprintf("%.1f", float64(weight)*100.0/float64(serialWeight))))
|
||||||
desc, _ := d.GetDescendants(rootIds[k])
|
|
||||||
|
|
||||||
for _, i := range desc {
|
|
||||||
ids = append(ids, i.(int))
|
|
||||||
}
|
|
||||||
|
|
||||||
sort.Ints(ids)
|
|
||||||
out(fmt.Sprintf("(%v) %v", len(ids), strings.Join(makeStrs(ids), "->")))
|
|
||||||
|
|
||||||
if len(desc) > maxDesc {
|
|
||||||
maxDesc = len(desc)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
numTx := len(d.DAG.GetVertices())
|
|
||||||
out(fmt.Sprintf("max chain length: %v of %v (%v%%)", maxDesc+1, numTx,
|
|
||||||
fmt.Sprintf("%.1f", float64(maxDesc+1)*100.0/float64(numTx))))
|
|
||||||
out(fmt.Sprintf("max dep count: %v of %v (%v%%)", maxDeps, totalDeps,
|
|
||||||
fmt.Sprintf("%.1f", float64(maxDeps)*100.0/float64(totalDeps))))
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -123,7 +123,7 @@ func (pq *SafePriorityQueue) Len() int {
|
||||||
|
|
||||||
type ParallelExecutionResult struct {
|
type ParallelExecutionResult struct {
|
||||||
TxIO *TxnInputOutput
|
TxIO *TxnInputOutput
|
||||||
Stats *[][]uint64
|
Stats *map[int]ExecutionStat
|
||||||
Deps *DAG
|
Deps *DAG
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -133,8 +133,9 @@ const numSpeculativeProcs = 8
|
||||||
type ParallelExecutor struct {
|
type ParallelExecutor struct {
|
||||||
tasks []ExecTask
|
tasks []ExecTask
|
||||||
|
|
||||||
// Stores the execution statistics for each task
|
// Stores the execution statistics for the last incarnation of each task
|
||||||
stats [][]uint64
|
stats map[int]ExecutionStat
|
||||||
|
|
||||||
statsMutex sync.Mutex
|
statsMutex sync.Mutex
|
||||||
|
|
||||||
// Channel for tasks that should be prioritized
|
// Channel for tasks that should be prioritized
|
||||||
|
|
@ -202,12 +203,20 @@ type ParallelExecutor struct {
|
||||||
workerWg sync.WaitGroup
|
workerWg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type ExecutionStat struct {
|
||||||
|
TxIdx int
|
||||||
|
Incarnation int
|
||||||
|
Start uint64
|
||||||
|
End uint64
|
||||||
|
Worker int
|
||||||
|
}
|
||||||
|
|
||||||
func NewParallelExecutor(tasks []ExecTask, profile bool) *ParallelExecutor {
|
func NewParallelExecutor(tasks []ExecTask, profile bool) *ParallelExecutor {
|
||||||
numTasks := len(tasks)
|
numTasks := len(tasks)
|
||||||
|
|
||||||
pe := &ParallelExecutor{
|
pe := &ParallelExecutor{
|
||||||
tasks: tasks,
|
tasks: tasks,
|
||||||
stats: make([][]uint64, numTasks),
|
stats: make(map[int]ExecutionStat, numTasks),
|
||||||
chTasks: make(chan ExecVersionView, numTasks),
|
chTasks: make(chan ExecVersionView, numTasks),
|
||||||
chSpeculativeTasks: make(chan struct{}, numTasks),
|
chSpeculativeTasks: make(chan struct{}, numTasks),
|
||||||
chSettle: make(chan int, numTasks),
|
chSettle: make(chan int, numTasks),
|
||||||
|
|
@ -272,10 +281,14 @@ func (pe *ParallelExecutor) Prepare() {
|
||||||
if pe.profile {
|
if pe.profile {
|
||||||
end := time.Since(pe.begin)
|
end := time.Since(pe.begin)
|
||||||
|
|
||||||
stat := []uint64{uint64(res.ver.TxnIndex), uint64(res.ver.Incarnation), uint64(start), uint64(end), uint64(procNum)}
|
|
||||||
|
|
||||||
pe.statsMutex.Lock()
|
pe.statsMutex.Lock()
|
||||||
pe.stats = append(pe.stats, stat)
|
pe.stats[res.ver.TxnIndex] = ExecutionStat{
|
||||||
|
TxIdx: res.ver.TxnIndex,
|
||||||
|
Incarnation: res.ver.Incarnation,
|
||||||
|
Start: uint64(start),
|
||||||
|
End: uint64(end),
|
||||||
|
Worker: procNum,
|
||||||
|
}
|
||||||
pe.statsMutex.Unlock()
|
pe.statsMutex.Unlock()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -348,8 +348,14 @@ func checkNoDroppedTx(pe *ParallelExecutor) error {
|
||||||
func runParallel(t *testing.T, tasks []ExecTask, validation PropertyCheck) time.Duration {
|
func runParallel(t *testing.T, tasks []ExecTask, validation PropertyCheck) time.Duration {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
profile := false
|
||||||
|
|
||||||
start := time.Now()
|
start := time.Now()
|
||||||
_, err := executeParallelWithCheck(tasks, false, validation)
|
result, err := executeParallelWithCheck(tasks, profile, validation)
|
||||||
|
|
||||||
|
if result.Deps != nil && profile {
|
||||||
|
result.Deps.Report(*result.Stats, func(str string) { fmt.Println(str) })
|
||||||
|
}
|
||||||
|
|
||||||
assert.NoError(t, err, "error occur during parallel execution")
|
assert.NoError(t, err, "error occur during parallel execution")
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,7 @@ package core
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/consensus"
|
"github.com/ethereum/go-ethereum/consensus"
|
||||||
|
|
@ -29,6 +30,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/core/vm"
|
"github.com/ethereum/go-ethereum/core/vm"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
"github.com/ethereum/go-ethereum/metrics"
|
||||||
"github.com/ethereum/go-ethereum/params"
|
"github.com/ethereum/go-ethereum/params"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -247,6 +249,8 @@ func (task *ExecutionTask) Settle() {
|
||||||
*task.allLogs = append(*task.allLogs, receipt.Logs...)
|
*task.allLogs = append(*task.allLogs, receipt.Logs...)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
var parallelizabilityTimer = metrics.NewRegisteredTimer("block/parallelizability", nil)
|
||||||
|
|
||||||
// Process processes the state changes according to the Ethereum rules by running
|
// Process processes the state changes according to the Ethereum rules by running
|
||||||
// the transaction messages using the statedb and applying any rewards to both
|
// the transaction messages using the statedb and applying any rewards to both
|
||||||
// the processor (coinbase) and any included uncles.
|
// the processor (coinbase) and any included uncles.
|
||||||
|
|
@ -314,7 +318,25 @@ func (p *ParallelStateProcessor) Process(block *types.Block, statedb *state.Stat
|
||||||
}
|
}
|
||||||
|
|
||||||
backupStateDB := statedb.Copy()
|
backupStateDB := statedb.Copy()
|
||||||
_, err := blockstm.ExecuteParallel(tasks, false)
|
|
||||||
|
profile := false
|
||||||
|
result, err := blockstm.ExecuteParallel(tasks, profile)
|
||||||
|
|
||||||
|
if err == nil && profile {
|
||||||
|
_, weight := result.Deps.LongestPath(*result.Stats)
|
||||||
|
|
||||||
|
serialWeight := uint64(0)
|
||||||
|
|
||||||
|
for i := 0; i < len(result.Deps.GetVertices()); i++ {
|
||||||
|
serialWeight += (*result.Stats)[i].End - (*result.Stats)[i].Start
|
||||||
|
}
|
||||||
|
|
||||||
|
parallelizabilityTimer.Update(time.Duration(serialWeight * 100 / weight))
|
||||||
|
|
||||||
|
log.Info("Parallelizability", "Average (%)", parallelizabilityTimer.Mean())
|
||||||
|
|
||||||
|
log.Info("Parallelizability", "Histogram (%)", parallelizabilityTimer.Percentiles([]float64{0.001, 0.01, 0.05, 0.1, 0.25, 0.5, 0.75, 0.9, 0.95, 0.99, 0.999, 0.9999}))
|
||||||
|
}
|
||||||
|
|
||||||
for _, task := range tasks {
|
for _, task := range tasks {
|
||||||
task := task.(*ExecutionTask)
|
task := task.(*ExecutionTask)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue