基于Goroutine实现任务依赖调度的发布订阅问题求助
修复Goroutine任务调度的数据竞争与依赖逻辑问题
以下是修复后的完整代码,解决了数据竞争问题并实现了正确的依赖任务调度逻辑:
package main import ( "fmt" "sync" ) type task struct { mu sync.Mutex isDone bool dependencies []*task subscribers []*task doneChan chan struct{} // 用空结构体实现轻量通知 numDependencies int taskName string } func (t *task) executeTask() { fmt.Printf("Task %s is getting executed...\n", t.taskName) // <-time.After(5 * time.Second) // 取消注释可模拟任务耗时 fmt.Printf("Task %s is done!! <-------\n", t.taskName) } // 安全标记任务完成并关闭通知通道 func (t *task) markDone() { t.mu.Lock() defer t.mu.Unlock() if !t.isDone { t.isDone = true close(t.doneChan) // 关闭通道作为任务完成的全局通知 } } // 安全添加订阅者 func (t *task) addSubscriber(sub *task) { t.mu.Lock() defer t.mu.Unlock() t.subscribers = append(t.subscribers, sub) } // 任务完成后通知所有订阅者减少依赖计数 func (t *task) notifySubscribers() { t.mu.Lock() // 复制订阅者切片避免持有锁期间遍历 subscribers := make([]*task, len(t.subscribers)) copy(subscribers, t.subscribers) t.mu.Unlock() for _, sub := range subscribers { fmt.Printf("Task %s notifying subscriber %s\n", t.taskName, sub.taskName) sub.decrementDependency() } } // 安全减少依赖计数,计数为0时触发任务执行 func (t *task) decrementDependency() { t.mu.Lock() defer t.mu.Unlock() t.numDependencies-- fmt.Printf("Task %s remaining dependencies: %d\n", t.taskName, t.numDependencies) if t.numDependencies == 0 { go t.run() // 依赖全部完成,启动当前任务 } } // 统一任务执行流程:执行 -> 标记完成 -> 通知订阅者 func (t *task) run() { t.executeTask() t.markDone() t.notifySubscribers() } func (t *task) setDependency(tasks []*task) { t.mu.Lock() defer t.mu.Unlock() t.dependencies = tasks t.numDependencies = len(tasks) } // 启动任务:无依赖直接执行,有依赖则监听所有依赖完成信号 func (t *task) start() { fmt.Printf("Starting Task %s\n", t.taskName) t.mu.Lock() // 复制依赖切片避免后续修改影响当前逻辑 deps := make([]*task, len(t.dependencies)) copy(deps, t.dependencies) initialDepCount := t.numDependencies t.mu.Unlock() if initialDepCount == 0 { go t.run() return } // 为每个依赖启动监听goroutine var wg sync.WaitGroup for _, dep := range deps { wg.Add(1) go func(d *task) { defer wg.Done() d.addSubscriber(t) <-d.doneChan // 等待依赖任务完成 t.decrementDependency() }(dep) } wg.Wait() fmt.Printf("Task %s is waiting for %d dependencies\n", t.taskName, initialDepCount) } func createTask(taskName string) *task { return &task{ isDone: false, taskName: taskName, dependencies: nil, subscribers: nil, numDependencies: 0, doneChan: make(chan struct{}), } } func main() { taskA := createTask("A") taskB := createTask("B") taskC := createTask("C") taskD := createTask("D") taskE := createTask("E") taskD.setDependency([]*task{taskA, taskB}) taskE.setDependency([]*task{taskC, taskD}) allTasks := []*task{taskA, taskB, taskC, taskD, taskE} var wg sync.WaitGroup for _, t := range allTasks { wg.Add(1) go func(t *task) { defer wg.Done() t.start() <-t.doneChan // 等待当前任务真正完成 }(t) } wg.Wait() fmt.Println("All tasks completed!") }
核心修复与优化说明
解决数据竞争
- 为
task结构体添加sync.Mutex,所有读写共享字段(isDone、numDependencies、subscribers)的操作都加锁保护。 - 遍历共享切片前先复制副本,避免持有锁期间执行耗时操作,降低锁竞争。
- 为
重构依赖通知逻辑
- 使用
chan struct{}替代chan bool,实现轻量的一次性完成通知(关闭通道是Go中标准的全局通知方式)。 - 每个依赖任务完成后,订阅者通过
decrementDependency安全减少计数,计数归零时自动启动当前任务,完全符合依赖调度逻辑。
- 使用
简化任务执行流程
- 移除原代码中冗余的
updateDependency、markCompleted等复杂逻辑,合并为run方法统一处理任务执行、标记完成、通知订阅者的全流程。 - 修复原代码中
trackDependency无差别启动goroutine导致重复减计数的问题,现在每个依赖对应一个监听goroutine,仅触发一次计数减少。
- 移除原代码中冗余的
完善主流程等待
- 主goroutine中每个任务的启动协程会等待
task.doneChan关闭,确保所有任务真正执行完成后再结束程序,避免提前退出。
- 主goroutine中每个任务的启动协程会等待
预期输出示例
Starting Task D Starting Task B Starting Task C Starting Task E Starting Task A Task B is getting executed... Task B is done!! <------- Task B notifying subscriber D Task D remaining dependencies: 1 Task C is getting executed... Task C is done!! <------- Task C notifying subscriber E Task E remaining dependencies: 1 Task A is getting executed... Task A is done!! <------- Task A notifying subscriber D Task D remaining dependencies: 0 Task D is getting executed... Task D is done!! <------- Task D notifying subscriber E Task E remaining dependencies: 0 Task E is getting executed... Task E is done!! <------- All tasks completed!
内容的提问来源于stack exchange,提问作者Pygirl
相关产品推荐
相关产品推荐

