|
|
@ -12,9 +12,9 @@
|
|
|
|
#define mpi_tag = 2008
|
|
|
|
#define mpi_tag = 2008
|
|
|
|
|
|
|
|
|
|
|
|
namespace paddle {
|
|
|
|
namespace paddle {
|
|
|
|
namespace operators {
|
|
|
|
namespace operators {
|
|
|
|
namespace detail {
|
|
|
|
namespace detail {
|
|
|
|
MPIUtils::MPIUtils(const std::string& worker_name) {
|
|
|
|
MPIUtils::MPIUtils(const std::string &worker_name) {
|
|
|
|
InitMPI();
|
|
|
|
InitMPI();
|
|
|
|
|
|
|
|
|
|
|
|
int rank = 0, size = 1;
|
|
|
|
int rank = 0, size = 1;
|
|
|
@ -29,9 +29,9 @@ MPIUtils::MPIUtils(const std::string& worker_name) {
|
|
|
|
for (int i = 0; i < number_of_procs; i++) {
|
|
|
|
for (int i = 0; i < number_of_procs; i++) {
|
|
|
|
name_to_id_[std::string(&worker_names[i * 128])] = i;
|
|
|
|
name_to_id_[std::string(&worker_names[i * 128])] = i;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void MPIUtils::InitMPI() {
|
|
|
|
void MPIUtils::InitMPI() {
|
|
|
|
int flag = 0;
|
|
|
|
int flag = 0;
|
|
|
|
MPI_CHECK(MPI_Initialized(&flag));
|
|
|
|
MPI_CHECK(MPI_Initialized(&flag));
|
|
|
|
|
|
|
|
|
|
|
@ -42,51 +42,49 @@ void MPIUtils::InitMPI() {
|
|
|
|
MPI_Init(0, 0);
|
|
|
|
MPI_Init(0, 0);
|
|
|
|
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
|
|
|
|
MPI_Comm_rank(MPI_COMM_WORLD, &rank);
|
|
|
|
MPI_Comm_size(MPI_COMM_WORLD, &size);
|
|
|
|
MPI_Comm_size(MPI_COMM_WORLD, &size);
|
|
|
|
MPI_Get_processor_name(host_name, &len)
|
|
|
|
MPI_Get_processor_name(host_name, &len);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
};
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
MPIIsend::MPIIsend(int dst, const char* req) {
|
|
|
|
|
|
|
|
done1 = 0;
|
|
|
|
MPISend::MPISend(const Meta &meta) {
|
|
|
|
done2 = 0;
|
|
|
|
done1_ = 1;
|
|
|
|
length = strlen(req);
|
|
|
|
done2_ = 0;
|
|
|
|
req = req;
|
|
|
|
this->meta = meta;
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
MPIIsend::Send() {
|
|
|
|
|
|
|
|
MPI_Isend(&req, length, MPI_CHAR, dst, mpi_tag, MPI_COMM_WORLD,
|
|
|
|
|
|
|
|
&msg1_);
|
|
|
|
|
|
|
|
MPI_Test(&msg1_, &done1_, MPI_STATUS_IGNORE)
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
bool MPIIsend::IsFinished() {
|
|
|
|
|
|
|
|
MPI_Status status;
|
|
|
|
|
|
|
|
if (!done1_) MPI_Test(&msg1_, &done1_, &status);
|
|
|
|
|
|
|
|
return done1;
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
MPIIsend::~MPIIsend(){
|
|
|
|
MPISend::Send() {
|
|
|
|
MPI_Wait(&msg1_, MPI_STATUS_IGNORE);
|
|
|
|
MPI_Send(&meta.request, meta.count, meta.datatype, meta.dst, meta.tag,
|
|
|
|
MPI_Free_mem(req);
|
|
|
|
MPI_COMM_WORLD);
|
|
|
|
}
|
|
|
|
done2_ = 1;
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
MPIIrecv::MPIIrecv(){
|
|
|
|
bool MPISend::IsReady() {
|
|
|
|
|
|
|
|
return true;
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
bool MPISend::IsFinished() { return done1_ && done2_; }
|
|
|
|
|
|
|
|
|
|
|
|
MPIIrecv::Recv(){
|
|
|
|
MPISend::~MPISend() { MPI_Free_mem(meta); }
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
MPIIrecv::IsFinished(){
|
|
|
|
MPIRecv::MPIRecv(const Meta &meta) {
|
|
|
|
|
|
|
|
this->meta = meta;
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
MPIRecv::Recv() {}
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
bool MPIRecv::IsReady() {
|
|
|
|
|
|
|
|
return true;
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
MPIIrecv::~MPIIrecv(){
|
|
|
|
MPIRecv::IsFinished() {}
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
MPIRecv::~MPIRecv() {
|
|
|
|
|
|
|
|
MPI_Free_mem(meta);
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
} // namespace detail
|
|
|
|
} // namespace detail
|
|
|
|
|
|
|
|
|
|
|
|
} // namespace operators
|
|
|
|
} // namespace operators
|
|
|
|
} // namespace paddle
|
|
|
|
} // namespace paddle
|