Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
123 changes: 26 additions & 97 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -176,110 +176,39 @@ brown dog lazy fox the quick brown dog the lazy fox`
}
```

## Example: Parallel Merge Sort

A real-world example using go-taskflow to parallelize merge sort across multiple chunks:

```go
package main

import (
"fmt"
"log"
"math/rand"
"os"
"slices"
"strconv"
"sync"

gtf "github.com/noneback/go-taskflow"
)

func mergeInto(dest, src []int) []int {
size := len(dest) + len(src)
tmp := make([]int, 0, size)
i, j := 0, 0
for i < len(dest) && j < len(src) {
if dest[i] < src[j] {
tmp = append(tmp, dest[i])
i++
} else {
tmp = append(tmp, src[j])
j++
}
}
if i < len(dest) {
tmp = append(tmp, dest[i:]...)
} else {
tmp = append(tmp, src[j:]...)
}
return tmp
}

func main() {
chunks := 100
chunkSize := 1000
randomArr := make([][]int, chunks)
sortedArr := make([]int, 0, chunks*chunkSize)
mutex := &sync.Mutex{}

for i := 0; i < chunks; i++ {
for j := 0; j < chunkSize; j++ {
randomArr[i] = append(randomArr[i], rand.Int())
}
}

sortTasks := make([]*gtf.Task, chunks)
tf := gtf.NewTaskFlow("merge sort")
done := tf.NewTask("done", func() {
if !slices.IsSorted(sortedArr) {
log.Fatal("sorting failed")
}
fmt.Println("sorted successfully")
})

for i := 0; i < chunks; i++ {
idx := i
sortTasks[idx] = tf.NewTask("sort_"+strconv.Itoa(idx), func() {
arr := randomArr[idx]
slices.Sort(arr)
mutex.Lock()
defer mutex.Unlock()
sortedArr = mergeInto(sortedArr, arr)
})
}
done.Succeed(sortTasks...)

executor := gtf.NewExecutor(1000, gtf.WithProfiler())
executor.Run(tf).Wait()

if err := tf.Dump(os.Stdout); err != nil {
log.Fatal(err)
}
if err := executor.Profile(os.Stdout); err != nil {
log.Fatal(err)
}
}
```

For more examples, visit the [examples directory](https://github.com/noneback/go-taskflow/tree/main/examples).

## Benchmark

The following benchmark provides a rough estimate of performance. Note that most realistic workloads are I/O-bound, and their performance cannot be accurately reflected by these results. For CPU-intensive tasks, consider using [taskflow-cpp](https://github.com/taskflow/taskflow).
The following benchmarks provide a rough estimate of pure scheduling overhead using empty task functions. Note that most realistic workloads are I/O-bound, and their performance cannot be accurately reflected by these results. For CPU-intensive tasks, consider using [taskflow-cpp](https://github.com/taskflow/taskflow).

```plaintext
$ go test -bench=. -benchmem
goos: linux
goarch: amd64
$ go test -bench=. -benchmem ./benchmark/
goos: darwin
goarch: arm64
pkg: github.com/noneback/go-taskflow/benchmark
cpu: Intel(R) Xeon(R) Platinum 8269CY CPU @ 2.50GHz
BenchmarkC32-4 23282 51891 ns/op 7295 B/op 227 allocs/op
BenchmarkS32-4 7047 160199 ns/op 6907 B/op 255 allocs/op
BenchmarkC6-4 66397 18289 ns/op 1296 B/op 47 allocs/op
BenchmarkC8x8-4 7946 143474 ns/op 16914 B/op 504 allocs/op
PASS
ok github.com/noneback/go-taskflow/benchmark 5.606s
cpu: Apple M4 Pro
BenchmarkConcurrent/N8-12 217042 5349 ns/op 1781 B/op 55 allocs/op
BenchmarkConcurrent/N32-12 47456 24439 ns/op 7566 B/op 213 allocs/op
BenchmarkConcurrent/N128-12 10000 116586 ns/op 32209 B/op 835 allocs/op
BenchmarkConcurrent/N512-12 2839 439930 ns/op 130337 B/op 3353 allocs/op
BenchmarkSerial/N8-12 126259 9339 ns/op 1905 B/op 63 allocs/op
BenchmarkSerial/N32-12 30313 39171 ns/op 7669 B/op 255 allocs/op
BenchmarkSerial/N128-12 7758 156781 ns/op 30725 B/op 1023 allocs/op
BenchmarkSerial/N512-12 1862 645739 ns/op 122952 B/op 4095 allocs/op
BenchmarkDiamond-12 181072 6662 ns/op 1441 B/op 47 allocs/op
BenchmarkDenseLayers/L4xW4-12 85122 13461 ns/op 4352 B/op 123 allocs/op
BenchmarkDenseLayers/L4xW8-12 42927 27127 ns/op 11764 B/op 270 allocs/op
BenchmarkDenseLayers/L8xW4-12 44412 27565 ns/op 8963 B/op 251 allocs/op
BenchmarkDenseLayers/L8xW8-12 20775 58088 ns/op 25071 B/op 556 allocs/op
BenchmarkSubflow-12 170228 6531 ns/op 1409 B/op 45 allocs/op
BenchmarkCondition-12 507645 2460 ns/op 704 B/op 23 allocs/op
BenchmarkConcurrencyScaling/C1-12 82400 14740 ns/op 15687 B/op 393 allocs/op
BenchmarkConcurrencyScaling/C12-12 23158 52602 ns/op 15700 B/op 422 allocs/op
BenchmarkConcurrencyScaling/C48-12 17500 68434 ns/op 15907 B/op 453 allocs/op
BenchmarkGraphBuild/N32-12 348994 3418 ns/op 6284 B/op 230 allocs/op
BenchmarkGraphBuild/N128-12 92428 13300 ns/op 24856 B/op 904 allocs/op
BenchmarkGraphBuild/N512-12 22143 55325 ns/op 101709 B/op 3850 allocs/op
```

## Understanding Conditional Tasks
Expand Down
189 changes: 138 additions & 51 deletions benchmark/benchmark_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,81 +2,168 @@ package benchmark

import (
"fmt"
"runtime"
"testing"

gotaskflow "github.com/noneback/go-taskflow"
)

var executor = gotaskflow.NewExecutor(6400)
// --- Topology scaling: measure scheduling overhead across graph shapes and sizes ---

func BenchmarkC32(b *testing.B) {
tf := gotaskflow.NewTaskFlow("G")
for i := 0; i < 32; i++ {
tf.NewTask(fmt.Sprintf("N%d", i), func() {})
func BenchmarkConcurrent(b *testing.B) {
for _, n := range []int{8, 32, 128, 512} {
b.Run(fmt.Sprintf("N%d", n), func(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("concurrent")
for i := 0; i < n; i++ {
tf.NewTask(fmt.Sprintf("T%d", i), func() {})
}
b.ResetTimer()
for i := 0; i < b.N; i++ {
exec.Run(tf).Wait()
}
})
}
}

func BenchmarkSerial(b *testing.B) {
for _, n := range []int{8, 32, 128, 512} {
b.Run(fmt.Sprintf("N%d", n), func(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("serial")
prev := tf.NewTask("T0", func() {})
for i := 1; i < n; i++ {
next := tf.NewTask(fmt.Sprintf("T%d", i), func() {})
prev.Precede(next)
prev = next
}
b.ResetTimer()
for i := 0; i < b.N; i++ {
exec.Run(tf).Wait()
}
})
}
}

func BenchmarkDiamond(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("diamond")
source := tf.NewTask("source", func() {})
left := tf.NewTask("left", func() {})
right := tf.NewTask("right", func() {})
leftMid := tf.NewTask("left_mid", func() {})
rightMid := tf.NewTask("right_mid", func() {})
sink := tf.NewTask("sink", func() {})
source.Precede(left, right)
left.Precede(leftMid)
right.Precede(rightMid)
sink.Succeed(leftMid, rightMid)
b.ResetTimer()
for i := 0; i < b.N; i++ {
executor.Run(tf).Wait()
exec.Run(tf).Wait()
}
}

func BenchmarkS32(b *testing.B) {
tf := gotaskflow.NewTaskFlow("G")
prev := tf.NewTask("N0", func() {})
for i := 1; i < 32; i++ {
next := tf.NewTask(fmt.Sprintf("N%d", i), func() {})
prev.Precede(next)
prev = next
func BenchmarkDenseLayers(b *testing.B) {
for _, layers := range []int{4, 8} {
for _, width := range []int{4, 8} {
b.Run(fmt.Sprintf("L%dxW%d", layers, width), func(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("dense_layers")
var curLayer, prevLayer []*gotaskflow.Task
for l := 0; l < layers; l++ {
for w := 0; w < width; w++ {
task := tf.NewTask(fmt.Sprintf("T%d_%d", l, w), func() {})
for _, p := range prevLayer {
p.Precede(task)
}
curLayer = append(curLayer, task)
}
prevLayer = curLayer
curLayer = nil
}
b.ResetTimer()
for i := 0; i < b.N; i++ {
exec.Run(tf).Wait()
}
})
}
}
}

// --- Feature benchmarks: paths unique to subflow and condition nodes ---

func BenchmarkSubflow(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("subflow")
setup := tf.NewTask("setup", func() {})
sub := tf.NewSubflow("sub", func(sf *gotaskflow.Subflow) {
s1 := sf.NewTask("sub_1", func() {})
s2 := sf.NewTask("sub_2", func() {})
s3 := sf.NewTask("sub_3", func() {})
s1.Precede(s3)
s2.Precede(s3)
})
teardown := tf.NewTask("teardown", func() {})
setup.Precede(sub)
sub.Precede(teardown)
b.ResetTimer()
for i := 0; i < b.N; i++ {
executor.Run(tf).Wait()
exec.Run(tf).Wait()
}
}

func BenchmarkC6(b *testing.B) {
tf := gotaskflow.NewTaskFlow("G")
n0 := tf.NewTask("N0", func() {})
n1 := tf.NewTask("N1", func() {})
n2 := tf.NewTask("N2", func() {})
n3 := tf.NewTask("N3", func() {})
n4 := tf.NewTask("N4", func() {})
n5 := tf.NewTask("N5", func() {})

n0.Precede(n1, n2)
n1.Precede(n3)
n2.Precede(n4)
n5.Succeed(n3, n4)

func BenchmarkCondition(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(runtime.NumCPU()))
tf := gotaskflow.NewTaskFlow("condition")
entry := tf.NewTask("entry", func() {})
cond := tf.NewCondition("cond", func() uint { return 0 })
branchA := tf.NewTask("branch_a", func() {})
branchB := tf.NewTask("branch_b", func() {})
entry.Precede(cond)
cond.Precede(branchA, branchB)
_ = branchA
_ = branchB
b.ResetTimer()
for i := 0; i < b.N; i++ {
executor.Run(tf).Wait()
exec.Run(tf).Wait()
}
}

func BenchmarkC8x8(b *testing.B) {
tf := gotaskflow.NewTaskFlow("G")

layersCount := 8
layerNodesCount := 8

var curLayer, upperLayer []*gotaskflow.Task

for i := 0; i < layersCount; i++ {
for j := 0; j < layerNodesCount; j++ {
task := tf.NewTask(fmt.Sprintf("N%d", i*layersCount+j), func() {})

for i := range upperLayer {
upperLayer[i].Precede(task)
// --- Concurrency scaling: fixed topology, varying executor concurrency ---

func BenchmarkConcurrencyScaling(b *testing.B) {
const taskCount = 64
numCPU := runtime.NumCPU()
for _, c := range []int{1, numCPU, numCPU * 4} {
b.Run(fmt.Sprintf("C%d", c), func(b *testing.B) {
exec := gotaskflow.NewExecutor(uint(c))
tf := gotaskflow.NewTaskFlow("conc_scaling")
for i := 0; i < taskCount; i++ {
tf.NewTask(fmt.Sprintf("T%d", i), func() {})
}

curLayer = append(curLayer, task)
}

upperLayer = curLayer
curLayer = []*gotaskflow.Task{}
b.ResetTimer()
for i := 0; i < b.N; i++ {
exec.Run(tf).Wait()
}
})
}
}

for i := 0; i < b.N; i++ {
executor.Run(tf).Wait()
// --- Graph construction: allocation cost of building DAGs ---

func BenchmarkGraphBuild(b *testing.B) {
for _, n := range []int{32, 128, 512} {
b.Run(fmt.Sprintf("N%d", n), func(b *testing.B) {
for i := 0; i < b.N; i++ {
tf := gotaskflow.NewTaskFlow("build")
prev := tf.NewTask("T0", func() {})
for j := 1; j < n; j++ {
next := tf.NewTask(fmt.Sprintf("T%d", j), func() {})
prev.Precede(next)
prev = next
}
}
})
}
}
Loading