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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 10:44:57