()
| 170 | } |
| 171 | |
| 172 | func (ms *Master) distributeMapTask() { |
| 173 | log.Trace("[Master] Start Map task") |
| 174 | |
| 175 | taskStates := sync.Map{} |
| 176 | |
| 177 | finish := make(chan string, 100) |
| 178 | // crashChan := make(chan MapTaskInfo, 100) |
| 179 | workerID := 0 |
| 180 | workers, nWorkers := ms.availableWorkers(ms.totalWorkers) |
| 181 | for _, mapTask := range ms.MapTasks { |
| 182 | mapTask.setState(TASK_INPROGRESS) |
| 183 | go func(task MapTaskInfo, id int) { |
| 184 | taskStates.Store(workers[id].UUID, task) |
| 185 | done := Map(workers[id].IP, task.toRPC()) |
| 186 | if !done { |
| 187 | // task.setState(TASK_IDLE) |
| 188 | workers[id].SetState(WORKER_UNKNOWN) |
| 189 | ms.crashChan <- workers[id].UUID |
| 190 | // crashChan <- task |
| 191 | } else { |
| 192 | // task.setState(TASK_COMPLETED) |
| 193 | workers[id].SetState(WORKER_IDLE) |
| 194 | finish <- workers[id].UUID |
| 195 | } |
| 196 | }(mapTask, workerID) |
| 197 | workerID = (workerID + 1) % nWorkers |
| 198 | } |
| 199 | |
| 200 | count := 0 |
| 201 | LOOP: |
| 202 | for { |
| 203 | select { |
| 204 | case <-finish: |
| 205 | count += 1 |
| 206 | if count == len(ms.MapTasks) { |
| 207 | close(finish) |
| 208 | break LOOP |
| 209 | } |
| 210 | case crashUUID := <-ms.crashChan: |
| 211 | // ms.waitForIDLEWorkers() |
| 212 | workers, _ := ms.availableWorkers(1) |
| 213 | log.Info("[Master] Re-execute Map Task from ", crashUUID, " To ", workers[0].UUID) |
| 214 | if crashUUID == workers[0].UUID { |
| 215 | log.Panic("Self loop") |
| 216 | } |
| 217 | reExecuteTask, ok := taskStates.Load(crashUUID) |
| 218 | if !ok { |
| 219 | log.Panic("Load task states from crashUUID fail") |
| 220 | } |
| 221 | go func(task MapTaskInfo, id int) { |
| 222 | // taskStates.Delete(crashUUID) |
| 223 | taskStates.Store(workers[id].UUID, task) |
| 224 | done := Map(workers[id].IP, task.toRPC()) |
| 225 | if !done { |
| 226 | task.setState(TASK_IDLE) |
| 227 | workers[id].SetState(WORKER_UNKNOWN) |
| 228 | ms.crashChan <- workers[id].UUID |
| 229 | } else { |
no test coverage detected