| 326 | |
| 327 | template<typename Dtype> |
| 328 | void NCCL<Dtype>::Run(const vector<int>& gpus, const char* restore) { |
| 329 | boost::barrier barrier(static_cast<int>(gpus.size())); |
| 330 | vector<NCCL<Dtype>*> nccls(gpus.size()); |
| 331 | // Create workers |
| 332 | vector<shared_ptr<Worker<Dtype> > > workers(gpus.size()); |
| 333 | for (int i = 1; i < gpus.size(); ++i) { |
| 334 | CUDA_CHECK(cudaSetDevice(gpus[i])); |
| 335 | Caffe::set_solver_rank(i); |
| 336 | Worker<Dtype>* w = new Worker<Dtype>(solver_, gpus[i], &barrier, |
| 337 | &nccls, restore); |
| 338 | w->StartInternalThread(); |
| 339 | workers[i].reset(w); |
| 340 | } |
| 341 | CUDA_CHECK(cudaSetDevice(gpus[0])); |
| 342 | Caffe::set_solver_rank(0); |
| 343 | barrier_ = &barrier; |
| 344 | solver_->add_callback(this); |
| 345 | if (solver_->param().layer_wise_reduce()) { |
| 346 | solver_->net()->add_after_backward(this); |
| 347 | } |
| 348 | nccls[0] = this; |
| 349 | // Wait for workers |
| 350 | barrier.wait(); |
| 351 | // Init NCCL |
| 352 | InitSingleProcess(&nccls); |
| 353 | barrier.wait(); |
| 354 | // Run first solver on current thread |
| 355 | Broadcast(); |
| 356 | solver_->Solve(); |
| 357 | barrier.wait(); // Hangs without it when running tests |
| 358 | // Wait for shutdown |
| 359 | for (int i = 1; i < gpus.size(); ++i) { |
| 360 | workers[i]->StopInternalThread(); |
| 361 | } |
| 362 | } |
| 363 | |
| 364 | INSTANTIATE_CLASS(Params); |
| 365 | INSTANTIATE_CLASS(GPUParams); |