You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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!")
}

核心修复与优化说明

  1. 解决数据竞争

    • 为task结构体添加sync.Mutex,所有读写共享字段(isDone、numDependencies、subscribers)的操作都加锁保护。
    • 遍历共享切片前先复制副本,避免持有锁期间执行耗时操作,降低锁竞争。
  2. 重构依赖通知逻辑

    • 使用chan struct{}替代chan bool,实现轻量的一次性完成通知(关闭通道是Go中标准的全局通知方式)。
    • 每个依赖任务完成后,订阅者通过decrementDependency安全减少计数,计数归零时自动启动当前任务,完全符合依赖调度逻辑。
  3. 简化任务执行流程

    • 移除原代码中冗余的updateDependency、markCompleted等复杂逻辑,合并为run方法统一处理任务执行、标记完成、通知订阅者的全流程。
    • 修复原代码中trackDependency无差别启动goroutine导致重复减计数的问题,现在每个依赖对应一个监听goroutine,仅触发一次计数减少。
  4. 完善主流程等待

    • 主goroutine中每个任务的启动协程会等待task.doneChan关闭,确保所有任务真正执行完成后再结束程序,避免提前退出。

预期输出示例

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 16:14:58