how to process async go function values in go

Viewed 683

GOAL

My goal is to find values in a matrix and send them to another async function to process

TRY

Create one go function for each row of a matrix to speed up the code and another go function to handle the matches - further go function spawning is not done yet.

EXPECTED:

I expected the go functions to collect data and send it through a channel to my collector, which continues processing (and might spawn new go functions)

RESULT:

My tries either ended in a deadlock or incomplete code execution (early return, so "collector finished" never showed up in console)

so I created a waitgroup then I start my go functions which will search for matches in my matrix

I added comments in the code to show what I already tried

func main() {
    var wg sync.WaitGroup

    // open a channel which listens for Points (line&row in the matrix)
    pts := make(chan Point)
    for i := 1; i <= lines; i++ {
        wg.Add(1) // add one to wait for the go function below...
        // send out a gopher to find a match
        go gopher(pts, /* more params*/ i, &wg)
    }
    
    // launch a collector to continue processing the collected values
    // if I add a wg.Add(1) here it ends up in a deadlock, as the collector contains a endless for loop
    // I want to start processing the results as soon as possible - not create a buffer and wait for the wg to finish
    // INFO: 'collector' does not collect all values from a stream and puts it into an array or slice. it continues processing -> raise new go functions -> level2
    go collector(pts, job)
    wg.Wait()
}

// NOTE: each function MAY return values 0..n
func gopher(pts chan, /* more params */, wg *sync.WaitGroup)
    defer wg.Done()
    //fmt.Printf("findGopher %d starting\n", y)
    // search algorithm here
    // any matches will be returned
    pts <- Point{x,y}
}

// this function is executed async and does not terminate - this is ok
// as long as processing is not finished
func collector(pts chan Point, job *Job) {
    //defer wg.Done() // does not help here.
    for {
        fmt.Print(/*"we found a match @",*/ <-pts, " ")
        // how can I add a break here to quit the loop when wg.wait() finished?
    }
    fmt.Println("collector finished.")
}

type Point struct {
    X int
    Y int
}

NOTE: I tried to search on stackoverflow and on google, without success, probably because I do not know how to name the issue (read about Condition Variables, Mutex and Semaphore Barrier with WaitGroup

SOLUTION

is there a simpler solution?

finally I created two waitgroups

wgFind
wgCollect

AND i close the channel when wgFind continues.

var wgFind sync.WaitGroup
var wgCollect sync.WaitGroup
pts := make(chan Point)

    for i := 1; i <= lines; i++ {
        wgFind.Add(1)
        // send out a gopher to find gopher in the image ;-)
        go searchForGopher(pts, /*more args*/, i, &wgFind)
    }

    // dont forget to collect the Points and log them
    wgCollect.Add(1)
    go collector(pts, &wgCollect)
    wgFind.Wait()
    close(pts) // signal that ALL go Functions are done.
    wgCollect.Wait()
0 Answers
Related