Getting there.
This commit is contained in:
@@ -11,6 +11,7 @@ set(SOURCES
|
||||
global.cc
|
||||
mailbox.h
|
||||
memory_mapped_file.h
|
||||
offset_dense_matrix.h
|
||||
point_request_message.h
|
||||
table.h
|
||||
transform.h
|
||||
|
||||
@@ -20,6 +20,7 @@
|
||||
#include "core/table/memory_mapped_file.h"
|
||||
#include "core/tree/gen_metric_tree.h"
|
||||
#include "core/table/distributed_auction.h"
|
||||
#include "core/table/offset_dense_matrix.h"
|
||||
#include "core/tree/distributed_local_kmeans.h"
|
||||
|
||||
namespace core {
|
||||
@@ -63,11 +64,65 @@ class DistributedTable: public boost::noncopyable {
|
||||
const std::vector<TreeType *> &top_leaf_nodes,
|
||||
int leaf_node_assignment_index) {
|
||||
|
||||
const int neighbor_radius = 2;
|
||||
const int num_iterations = 10;
|
||||
|
||||
// Readjust the centroid.
|
||||
std::vector<int> point_assignments;
|
||||
int total_num_points_owned;
|
||||
core::tree::DistributedLocalKMeans local_kmeans;
|
||||
local_kmeans.Compute(
|
||||
table_outbox_group_comm, metric, *owned_table_,
|
||||
neighbor_radius, num_iterations,
|
||||
top_leaf_nodes[leaf_node_assignment_index]->bound().center(),
|
||||
2, 10);
|
||||
&total_num_points_owned, &point_assignments);
|
||||
|
||||
// Move the data across processes to get a new local table.
|
||||
TableType *new_local_table =
|
||||
core::table::global_m_file_->Construct<TableType>();
|
||||
new_local_table->Init(
|
||||
owned_table_->n_attributes(), total_num_points_owned);
|
||||
|
||||
// Left contributions.
|
||||
std::vector < core::table::OffsetDenseMatrix > left_contributions;
|
||||
std::vector < core::table::OffsetDenseMatrix > right_contributions;
|
||||
std::vector< boost::mpi::request > left_send_requests;
|
||||
std::vector< boost::mpi::request > right_send_requests;
|
||||
left_contributions.resize(
|
||||
std::min(table_outbox_group_comm.rank(), neighbor_radius));
|
||||
right_contributions.resize(
|
||||
std::min(
|
||||
table_outbox_group_comm.size() -
|
||||
table_outbox_group_comm.rank() - 1, neighbor_radius));
|
||||
for(unsigned int i = 1; i <= left_contributions.size(); i++) {
|
||||
|
||||
// Send and receive.
|
||||
left_send_requests_[i - 1] =
|
||||
comm.isend(comm.rank() - i, i, local_centroid_);
|
||||
left_receive_requests_[i - 1] =
|
||||
comm.irecv(
|
||||
comm.rank() - i, neighbor_radius + i, left_centroids_[i - 1]);
|
||||
}
|
||||
std::vector< core::table::OffsetDenseMatrix > right_contributions;
|
||||
for(unsigned int i = 1; i <= right_centroids_.size(); i++) {
|
||||
|
||||
// Send and receive.
|
||||
right_send_requests_[i - 1] =
|
||||
comm.isend(
|
||||
comm.rank() + i, neighbor_radius + i, local_centroid_);
|
||||
right_receive_requests_[i - 1] =
|
||||
comm.irecv(comm.rank() + i, i, right_centroids_[i - 1]);
|
||||
}
|
||||
|
||||
// Wait for all send/receive requests to be completed.
|
||||
boost::mpi::wait_all(
|
||||
left_send_requests_.begin(), left_send_requests_.end());
|
||||
boost::mpi::wait_all(
|
||||
left_receive_requests_.begin(), left_receive_requests_.end());
|
||||
boost::mpi::wait_all(
|
||||
right_send_requests_.begin(), right_send_requests_.end());
|
||||
boost::mpi::wait_all(
|
||||
right_receive_requests_.begin(), right_receive_requests_.end());
|
||||
}
|
||||
|
||||
int TakeLeafNodeOwnerShip_(
|
||||
|
||||
@@ -0,0 +1,101 @@
|
||||
/** @file offset_dense_matrix.h
|
||||
*
|
||||
* @author Dongryeol Lee (dongryel@cc.gatech.edu)
|
||||
*/
|
||||
|
||||
#ifndef CORE_TABLE_OFFSET_DENSE_MATRIX_H
|
||||
#define CORE_TABLE_OFFSET_DENSE_MATRIX_H
|
||||
|
||||
#include "core/table/dense_matrix.h"
|
||||
#include <boost/serialization/serialization.hpp>
|
||||
|
||||
namespace core {
|
||||
namespace table {
|
||||
class OffsetDenseMatrix {
|
||||
private:
|
||||
friend class boost::serialization::access;
|
||||
|
||||
double *ptr_;
|
||||
|
||||
std::vector<int> *assignment_indices_;
|
||||
|
||||
int filter_index_;
|
||||
|
||||
int n_attributes_;
|
||||
|
||||
int n_entries_;
|
||||
|
||||
public:
|
||||
|
||||
int n_entries() const {
|
||||
return n_entries_;
|
||||
}
|
||||
|
||||
OffsetDenseMatrix() {
|
||||
ptr_ = NULL;
|
||||
assignment_indices_ = NULL;
|
||||
filter_index_ = -1;
|
||||
n_attributes_ = -1;
|
||||
n_entries_ = -1;
|
||||
}
|
||||
|
||||
void Init(double *ptr_in, int n_attributes_in) {
|
||||
ptr_ = ptr_in;
|
||||
n_attributes_ = n_attributes_in;
|
||||
}
|
||||
|
||||
void set_filter_index(int filter_index_in) {
|
||||
filter_index_ = filter_index_in;
|
||||
}
|
||||
|
||||
void Init(
|
||||
core::table::DenseMatrix &mat_in,
|
||||
std::vector<int> &assignment_indices_in) {
|
||||
ptr_ = mat_in.ptr();
|
||||
n_attributes_ = mat_in.n_attributes();
|
||||
n_entries_ = mat_in.n_entries();
|
||||
assignment_indices_ = &assignment_indices_in;
|
||||
}
|
||||
|
||||
template<class Archive>
|
||||
void save(Archive &ar, const unsigned int version) const {
|
||||
|
||||
// First, save the number of doubles to be serialized.
|
||||
int num_doubles = 0;
|
||||
for(int i = 0; i < assignment_indices_->size(); i++) {
|
||||
if((*assignment_indices_)[i] == filter_index_) {
|
||||
num_doubles++;
|
||||
}
|
||||
}
|
||||
num_doubles *= n_attributes_;
|
||||
ar & num_doubles;
|
||||
|
||||
// Loop through and find out the columns to serialize.
|
||||
double *ptr_iter = ptr_;
|
||||
for(int i = 0; i < assignment_indices_->size();
|
||||
i++, ptr_iter += n_attributes_) {
|
||||
if((*assignment_indices_)[i] == filter_index_) {
|
||||
for(int j = 0; j < n_attributes_; j++) {
|
||||
ar & ptr_iter[j];
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
template<class Archive>
|
||||
void load(Archive &ar, const unsigned int version) {
|
||||
|
||||
// Load the number of points to be unfrozen.
|
||||
int num_doubles;
|
||||
ar & num_doubles;
|
||||
for(int i = 0; i < num_doubles; i++) {
|
||||
ar & ptr_[i];
|
||||
}
|
||||
n_entries_ = num_doubles / n_attributes_;
|
||||
}
|
||||
BOOST_SERIALIZATION_SPLIT_MEMBER()
|
||||
};
|
||||
};
|
||||
};
|
||||
|
||||
#endif
|
||||
@@ -246,7 +246,9 @@ class Table: public boost::noncopyable {
|
||||
|
||||
void Init(
|
||||
int num_dimensions_in, int num_points_in) {
|
||||
data_.Init(num_dimensions_in, num_points_in);
|
||||
if(num_dimensions_in > 0 && num_points_in > 0) {
|
||||
data_.Init(num_dimensions_in, num_points_in);
|
||||
}
|
||||
}
|
||||
|
||||
void Init(const std::string &file_name) {
|
||||
|
||||
+25
-14
@@ -101,14 +101,12 @@ class DistributedLocalKMeans {
|
||||
|
||||
std::vector< CentroidInfo > tmp_recv_from_right_;
|
||||
|
||||
std::vector<int> point_assignments_;
|
||||
|
||||
private:
|
||||
|
||||
template<typename TableType>
|
||||
void SynchronizeCentroids_(
|
||||
boost::mpi::communicator &comm, TableType &local_table_in,
|
||||
int neighbor_radius) {
|
||||
int neighbor_radius, const std::vector<int> &point_assignments) {
|
||||
|
||||
// Reset the contribution list.
|
||||
local_centroid_.Reset();
|
||||
@@ -119,19 +117,19 @@ class DistributedLocalKMeans {
|
||||
tmp_right_centroids_[i].Reset();
|
||||
}
|
||||
|
||||
for(unsigned int i = 0; i < point_assignments_.size(); i++) {
|
||||
for(unsigned int i = 0; i < point_assignments.size(); i++) {
|
||||
core::table::DensePoint point;
|
||||
local_table_in.get(i, &point);
|
||||
|
||||
if(point_assignments_[i] == comm.rank()) {
|
||||
if(point_assignments[i] == comm.rank()) {
|
||||
local_centroid_.Add(point);
|
||||
}
|
||||
else if(point_assignments_[i] < comm.rank()) {
|
||||
int destination_index = comm.rank() - point_assignments_[i] - 1;
|
||||
else if(point_assignments[i] < comm.rank()) {
|
||||
int destination_index = comm.rank() - point_assignments[i] - 1;
|
||||
tmp_left_centroids_[destination_index].Add(point);
|
||||
}
|
||||
else {
|
||||
int destination_index = point_assignments_[i] - comm.rank() - 1;
|
||||
int destination_index = point_assignments[i] - comm.rank() - 1;
|
||||
tmp_right_centroids_[destination_index].Add(point);
|
||||
}
|
||||
}
|
||||
@@ -187,8 +185,10 @@ class DistributedLocalKMeans {
|
||||
boost::mpi::communicator &comm,
|
||||
const core::metric_kernels::AbstractMetric &metric,
|
||||
TableType &local_table_in,
|
||||
const core::table::DensePoint &starting_centroid,
|
||||
int neighbor_radius, int num_outer_loop_iterations) {
|
||||
int neighbor_radius, int num_outer_loop_iterations,
|
||||
core::table::DensePoint &starting_centroid,
|
||||
int *total_num_points_owned,
|
||||
std::vector<int> *point_assignments_out) {
|
||||
|
||||
// Every process collects the local centers from the process ID
|
||||
// within the specified neighbor_radius.
|
||||
@@ -221,7 +221,7 @@ class DistributedLocalKMeans {
|
||||
}
|
||||
|
||||
// List of assignments for each point in the table.
|
||||
point_assignments_.resize(local_table_in.n_entries());
|
||||
point_assignments_out->resize(local_table_in.n_entries());
|
||||
|
||||
// The outer loop.
|
||||
for(
|
||||
@@ -289,21 +289,32 @@ class DistributedLocalKMeans {
|
||||
}
|
||||
|
||||
// Assign the point.
|
||||
point_assignments_[p] = min_index;
|
||||
(*point_assignments_out)[p] = min_index;
|
||||
|
||||
} // end of iterating over each point.
|
||||
|
||||
// Recompute the local centroid contribution and send to the
|
||||
// left and to the right.
|
||||
SynchronizeCentroids_(comm, local_table_in, neighbor_radius);
|
||||
SynchronizeCentroids_(
|
||||
comm, local_table_in, neighbor_radius, *point_assignments_out);
|
||||
|
||||
// Synchronize before continuing the outer iteration.
|
||||
comm.barrier();
|
||||
|
||||
} // end of outer iterations.
|
||||
|
||||
printf("Local ended up %d from %d\n", local_centroid_.num_points(),
|
||||
// Copy the final ending position and the total number of points
|
||||
// owned by this centroid.
|
||||
starting_centroid.CopyValues(local_centroid_.centroid());
|
||||
*total_num_points_owned = local_centroid_.num_points();
|
||||
|
||||
printf("%d: Local ended up %d from %d\n", comm.rank(),
|
||||
local_centroid_.num_points(),
|
||||
local_table_in.n_entries());
|
||||
for(unsigned int i = 0; i < point_assignments_out->size(); i++) {
|
||||
printf("%d ", (*point_assignments_out)[i]);
|
||||
}
|
||||
printf("\n");
|
||||
}
|
||||
};
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user