(files []string, nWorker int, nReduce int, addr string)
| 25 | } |
| 26 | |
| 27 | func StartMaster(files []string, nWorker int, nReduce int, addr string) { |
| 28 | // start gRPC server |
| 29 | listener, err := net.Listen("tcp", addr) |
| 30 | if err != nil { |
| 31 | log.Panic(err) |
| 32 | } |
| 33 | ms := NewMaster(nWorker, nReduce) |
| 34 | baseServer := grpc.NewServer() |
| 35 | rpc.RegisterMasterServer(baseServer, ms) |
| 36 | go baseServer.Serve(listener) |
| 37 | |
| 38 | log.Info("[Master] Master gRPC server start") |
| 39 | |
| 40 | // Check the worker is enough |
| 41 | ms.(*Master).waitForEnoughWorker() |
| 42 | go ms.(*Master).PeriodicHealthCheck() |
| 43 | |
| 44 | // Split input file (100,000 lines per chunk) |
| 45 | ms.(*Master).distributeWork(files) |
| 46 | |
| 47 | ms.(*Master).distributeMapTask() |
| 48 | |
| 49 | ms.(*Master).distributeReduceTask() |
| 50 | |
| 51 | ms.(*Master).endWorkers() |
| 52 | |
| 53 | baseServer.Stop() |
| 54 | } |
no test coverage detected