대량의 데이터를 나누어 처리할 때 데이터 개수를 워커마다 균등하게 배분해서 쓰고는 했는데,

데이터 당 처리시간이 불균질할 경우 어떤 워커는 바삐 일하고 어떤 워커는 노는 일이 발생하고는 했습니다.


그래서 남들은 어찌할까 찾아보려다가 일단 스스로 해보자는 생각에 짜보았습니다.

고루틴/채널의 특성을 이용해 콤퓨타가 알아서 노는 워커에게 일감을 건네주도록 하였습니다.



아주 허접한 시뮬레이션이지만 의도한 대로 굴러는 가네요.

막판까지 워커들이 골고루 일하면서 각자 처리한 데이터 숫자는 다릅니다.

참고로 데이터 처리시간은 데이터 생성시 랜덤하게 처리시간을 정하도록 하였습니다. (1~500ms)


재미있네요.

한 번 동시성 프로그래밍에 대해 제대로 공부해봐야 할까봐요.

하지만 다른 언어로 하면 코피 터지게 어렵겠죠...........




package main


import (

    "context"

    "fmt"

    "math/rand"

    "sync"

    "time"

)


const (

    DATA_NUM   = 500

    WORKER_NUM = 10

)


func main() {

    ctx, cancel := context.WithCancel(context.Background())

    ch := make(chan *Data)

    wg := new(sync.WaitGroup)

    wg.Add(WORKER_NUM)


    workers := newWorkers(WORKER_NUM)

    dataset := newDataSet(DATA_NUM)


    go workers.Work(ctx, wg, ch)

    go dataset.Send(cancel, ch)


    wg.Wait()

    workers.SayProcessed()

}


type Worker struct {

    ID        int

    Processed int

}


type Workers []*Worker


func newWorkers(num int) Workers {

    var workers Workers

    for i := 1; i < num+1; i++ {

        workers = append(workers, &Worker{

            ID:        i,

            Processed: 0,

        })

    }

    return workers

}


func (workers Workers) Work(ctx context.Context, wg *sync.WaitGroup, ch <-chan *Data) {

    for _, worker := range workers {

        go func(ctx context.Context, wg *sync.WaitGroup, ch <-chan *Data, worker *Worker) {

            defer wg.Done()

            for {

                select {

                case data := <-ch:

                    time.Sleep(data.Duration)

                    worker.Processed++

                    fmt.Printf("[Worker d] has just processed data d!\n", worker.ID, data.ID)

                case <-ctx.Done():

                    return

                }

            }

        }(ctx, wg, ch, worker)

    }

}


func (workers Workers) SayProcessed() {

    total := 0

    fmt.Println("----------------------------------------")

    for _, worker := range workers {

        fmt.Printf("[Worker d] processed %d data.\n", worker.ID, worker.Processed)

        total += worker.Processed

    }

    fmt.Println("----------------------------------------")

    fmt.Printf("Workers processed %d data.\n", total)

}


type Data struct {

    ID       int

    Duration time.Duration

}


type DataSet []*Data


func newDataSet(num int) DataSet {

    var dataSet DataSet

    rand.Seed(time.Now().UnixNano())

    for i := 1; i < num+1; i++ {

        dataSet = append(dataSet, &Data{

            ID:       i,

            Duration: time.Millisecond * time.Duration(rand.Intn(500)+1),

        })

    }

    return dataSet

}


func (dataSet DataSet) Send(cancel context.CancelFunc, ch chan<- *Data) {

    for _, data := range dataSet {

        ch <- data

    }

    cancel()

}