GoLang分布式系统任务分配异常:Worker无法获取有效任务
问题描述
我正在用GoLang结合hashcat构建分布式密码破解器,目前在调试协调者(Master)与Worker节点的任务分配逻辑,还没添加密码破解核心代码。手动创建了包含2个任务的任务队列,但Worker节点一直收到ID为0的空任务(已通过判断ID为0输出“无可用任务”)。以下是Coordinator.go与Worker.go的代码,请问任务分配逻辑哪里出问题了?
Coordinator.go
package main import ( "fmt" "net" "net/rpc" "sync" ) type Task struct { ID int HashToCrack string // Wordlist string } type Result struct { TaskID int Hash string Cracked bool Password string } type WorkerNode int type MasterNode struct { workerNodes []string taskQueue chan Task results chan Result mutex sync.Mutex } func (m *MasterNode) GetTask(_ string, task *Task) error { select { case t := <-m.taskQueue: *task = t default: // No tasks available at the moment *task = Task{} } return nil } // ReportResult receives and stores results from worker nodes. func (m *MasterNode) ReportResult(result Result, reply *bool) error { m.results <- result *reply = true return nil } func (m *MasterNode) RegisterWorkerNode(address string, reply *bool) error { m.mutex.Lock() defer m.mutex.Unlock() m.workerNodes = append(m.workerNodes, address) *reply = true return nil } // DistributeTasks distributes tasks to worker nodes. func (m *MasterNode) DistributeTasks(_ Task, _ *bool) error { m.mutex.Lock() defer m.mutex.Unlock() // Distribute tasks to worker nodes for _, worker := range m.workerNodes { select { case task := <-m.taskQueue: go func(worker string, task Task) { client, err := rpc.Dial("tcp", worker) if err != nil { fmt.Printf("Error connecting to worker node %s: %v\n", worker, err) return } defer client.Close() var workerReply bool err = client.Call("WorkerNode.CrackPassword", task, &workerReply) if err != nil { fmt.Printf("Error calling worker node %s: %v\n", worker, err) return } }(worker, task) default: // No tasks available to distribute } } return nil } func main() { master := new(MasterNode) rpc.Register(master) // Start RPC server for master listener, err := net.Listen("tcp", ":1234") if err != nil { fmt.Println("Error starting RPC server:", err) return } defer listener.Close() fmt.Println("Master Node RPC server is listening on port 1234") // Create a task queue and results channel master.taskQueue = make(chan Task) master.results = make(chan Result) // Simulate distributing tasks to worker nodes go func() { tasks := []Task{ {ID: 1, HashToCrack: "hash1"}, {ID: 2, HashToCrack: "hash2"}, } for _, task := range tasks { master.taskQueue <- task } }() // Handle task distribution go func() { if master.taskQueue != nil { for task := range master.taskQueue { var reply bool err := master.DistributeTasks(task, &reply) if err != nil { fmt.Println("Error distributing tasks:", err) } } } }() // Process task results go func() { for result := range master.results { fmt.Printf("Task %d: Hash: %s, Cracked: %v, Password: %s\n", result.TaskID, result.Hash, result.Cracked, result.Password) } }() // Accept incoming RPC connections for { conn, err := listener.Accept() if err != nil { fmt.Println("Error accepting RPC connection:", err) continue } go rpc.ServeConn(conn) } }
Worker.go
package main import ( "fmt" "net/rpc" "time" ) type Task struct { ID int HashToCrack string // Wordlist string } type WorkerNode int type Result struct { TaskID int Hash string Cracked bool Password string } func (w *WorkerNode) CrackPassword(task Task, reply *bool) error { // Simulated password cracking logic (replace with your actual implementation) // For simplicity, we'll assume the password is "123456" for any hash. hash := task.HashToCrack password := "123456" // Simulate the time it takes to crack the password (for demonstration purposes) time.Sleep(2 * time.Second) // Send the result back to the master result := Result{ TaskID: task.ID, Hash: hash, Cracked: true, Password: password, // Replace with the actual cracked password } master, err := rpc.Dial("tcp", "127.0.0.1"+":1234") if err != nil { fmt.Println("Error connecting to master node:", err) return err } defer master.Close() var masterReply bool err = master.Call("MasterNode.ReportResult", result, &masterReply) if err != nil { fmt.Println("Error reporting result to master node:", err) return err } *reply = true return nil } func main() { // Register with the master node w := new(WorkerNode) master, err := rpc.Dial("tcp", "127.0.0.1"+":1234") if err != nil { fmt.Println("Error connecting to master node:", err) return } defer master.Close() var reply bool err = master.Call("MasterNode.RegisterWorkerNode", "127.0.0.1"+":1235", &reply) if err != nil { fmt.Println("Error registering with master node:", err) return } // Simulate requesting tasks and reporting results to the master for { var task Task err := master.Call("MasterNode.GetTask", "", &task) if err != nil { fmt.Println("Error getting task from master node:", err) return } // Check if the task is empty, indicating no tasks are available if task.ID == 0 { fmt.Println("No tasks available. Waiting for tasks...") time.Sleep(5 * time.Second) // Wait for tasks to become available continue } // Handle the task if it's not empty fmt.Printf("Worker: Received Task %d - Hash: %s\n", task.ID, task.HashToCrack) // Perform password cracking (calling the CrackPassword method) var crackReply bool err = w.CrackPassword(task, &crackReply) if err != nil { fmt.Println("Error performing password cracking:", err) return } } }
手动创建的任务队列
tasks := []Task{ {ID: 1, HashToCrack: "hash1"}, {ID: 2, HashToCrack: "hash2"}, }
问题分析与修复方案
你的任务分配逻辑存在两个核心问题:
1. 任务队列被重复消费,导致Worker无任务可拿
Master的main函数里启动了两个goroutine消费taskQueue:
- 第一个goroutine负责写入任务到队列
- 第二个goroutine(Handle task distribution)通过
for task := range master.taskQueue持续读取队列,调用DistributeTasks后,该方法又会再次从taskQueue读取任务尝试推送给Worker
这导致任务刚写入队列就被第二个goroutine读走,后续DistributeTasks读取队列时已经为空,而Worker调用GetTask时,队列里已经没有剩余任务,只能返回空任务。
修复:
删除Handle task distribution的goroutine,同时删除DistributeTasks方法及相关调用。因为当前Worker是主动拉取任务模式,不需要Master主动推送任务,只需要保留GetTask方法供Worker拉取即可。
2. 无缓冲队列可能导致写入阻塞
当前master.taskQueue = make(chan Task)是无缓冲队列,当写入任务时如果没有消费者读取,会导致写入阻塞。建议改为带缓冲的队列,避免阻塞问题。
修复:
初始化队列时指定缓冲大小:
master.taskQueue = make(chan Task, 10) // 缓冲大小可根据实际需求调整
修改后的Master核心代码示例
修改后的main函数:
func main() { master := new(MasterNode) rpc.Register(master) // Start RPC server for master listener, err := net.Listen("tcp", ":1234") if err != nil { fmt.Println("Error starting RPC server:", err) return } defer listener.Close() fmt.Println("Master Node RPC server is listening on port 1234") // Create a buffered task queue and results channel master.taskQueue = make(chan Task, 10) master.results = make(chan Result) // Simulate adding tasks to queue go func() { tasks := []Task{ {ID: 1, HashToCrack: "hash1"}, {ID: 2, HashToCrack: "hash2"}, } for _, task := range tasks { master.taskQueue <- task } }() // Process task results go func() { for result := range master.results { fmt.Printf("Task %d: Hash: %s, Cracked: %v, Password: %s\n", result.TaskID, result.Hash, result.Cracked, result.Password) } }() // Accept incoming RPC connections for { conn, err := listener.Accept() if err != nil { fmt.Println("Error accepting RPC connection:", err) continue } go rpc.ServeConn(conn) } }
修改后,Worker调用GetTask时就能正确从队列中获取到任务,不会再收到ID为0的空任务。
内容的提问来源于stack exchange,提问作者Ehtesham
相关产品推荐
相关产品推荐

