Je rencontre parfois ce problème après l'avoir exécuté encore et encore. Je sais que c'est lié aux compteurs. Il renvoie cette erreur lorsque la méthode done() de sync.waitgroup est appelée plus de fois que la méthode add() n'est appelée.
Comment résoudre ce problème ?
Mon code crée des lots de taille 4 et effectue des traitements sur chaque lot, mais j'ai du mal à résoudre cette panique.
package main import ( "fmt" "sync" ) func main() { // create input channel input := make(chan int) // create wait group var wg sync.waitgroup // start batcher goroutine wg.add(1) go batcher(input, &wg) // send input values to the batcher for i := 1; i <= 10; i++ { input <- i } // close input channel close(input) // wait for batcher goroutine to finish wg.wait() } func batcher(input chan int, wg *sync.waitgroup) { // create batch channel with buffer of size 4 batch := make(chan int, 4) // create channel to synchronize worker goroutines done := make(chan bool) // create wait group for worker goroutines var workerwg sync.waitgroup // start worker goroutines for i := 0; i < 4; i++ { workerwg.add(1) go worker(batch, &workerwg, done) } // read input values and send to batch for value := range input { batch <- value if len(batch) == 4 { // wait for worker goroutines to finish processing batch workerwg.wait() // send batch to worker goroutines for i := 0; i < 4; i++ { workerwg.add(1) go sendbatch(batch, &workerwg, done) } } } // wait for worker goroutines to finish processing remaining batch workerwg.wait() // close done channel to notify that all batches have been processed close(done) wg.done() } func sendbatch(batch chan int, workerwg *sync.waitgroup, done chan bool) { // process batch for value := range batch { fmt.println("processing value:", value) } // notify worker goroutines that batch has been processed workerwg.done() select { case done <- true: default: // done channel has been closed } } func worker(batch chan int, workerwg *sync.waitgroup, done chan bool) { // process batches received from batch channel for batch := range batch { // process batch fmt.println("processing batch:", batch) workerwg.done() } // notify batcher goroutine that worker goroutine has finished select { case done <- true: default: // done channel has been closed } }
Code de base pour l'écriture de programmes batch :
package main import ( "fmt" "sync" ) func main() { input := make(chan int) output := make(chan []int) var wg sync.waitgroup wg.add(2) // start the batcher goroutine go func() { batch := []int{} for value := range input { batch = append(batch, value) if len(batch) == 4 { output <- batch batch = []int{} } } if len(batch) > 0 { output <- batch } close(output) wg.done() }() // start the worker goroutine go func() { for batch := range output { sum := 0 for _, value := range batch { sum += value } fmt.printf("sum of batch %v: %d\n", batch, sum) } wg.done() }() // send input values to the batcher for _, v := range []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10} { input <- v } close(input) // wait for both goroutines to finish wg.wait() }
Sum of batch [1 2 3 4]: 10 Sum of batch [5 6 7 8]: 26 Sum of batch [9 10]: 19
La première conception est un peu complexe, je vais essayer de développer cette conception de base.
Selon ce code :
for i := 0; i < 4; i++ { workerwg.add(1) go worker(batch, &workerwg, done) }
Je pense que workerwg.done()
devrait être déplacé en dehors de la boucle :
func worker(batch chan int, workerWg *sync.WaitGroup, done chan bool) { + defer workerWg.Done() // process batches received from batch channel for batch := range batch { // process batch fmt.Println("Processing batch:", batch) - workerWg.Done() } // notify batcher goroutine that worker goroutine has finished select { case done <- true: default: // done channel has been closed } }
Mais batch
n'est pas fermé dans la démo. Donc en fait, la goroutine fonctionnera indéfiniment jusqu'à la fin du programme.
Je ne sais pas s'il y a d'autres questions. La conception est trop compliquée. Le code complexe est difficile à comprendre et sujet aux erreurs. Pensez à le repenser.
Ce qui précède est le contenu détaillé de. pour plus d'informations, suivez d'autres articles connexes sur le site Web de PHP en chinois!