السلام عليكم ورحمة الله وبركاته
ارجو من مبرمجي وخبراء الـ c# المحترمون في هذا المنتدى مساعدتي وذلك بتحويل البرامج التالية من لغة Java الى C#
علما اني استخدمت برامج التحويل الجاهز للكود لكنها لم تنفع
ولكم جزيل الشكر والتقدير
1- BasicKMeans
package kmeans;
import java.util.*;
/**
* Basic implementation of K-means clustering. Since it's a Runnable, it's
* designed to be executed by a dedicated thread, but that thread
* does not create any other threads to divide up the work.
*/
public class BasicKMeans implements KMeans {
// Temporary clusters used during the clustering process. Converted to
// an array of the simpler class Cluster at the conclusion.
private ProtoCluster[] mProtoClusters;
// Cache of coordinate-to-cluster distances. Number of entries =
// number of clusters X number of coordinates.
private double[][] mDistanceCache;
// Used in makeAssignments() to figure out how many moves are made
// during each iteration -- the cluster assignment for coordinate n is
// found in mClusterAssignments[n] where the N coordinates are numbered
// 0 ... (N-1)
private int[] mClusterAssignments;
// 2D array holding the coordinates to be clustered.
private double[][] mCoordinates;
// The desired number of clusters and maximum number
// of iterations.
private int mK, mMaxIterations;
// Seed for the random number generator used to select
// coordinates for the initial cluster centers.
private long mRandomSeed;
// An array of Cluster objects: the output of k-means.
private Cluster[] mClusters;
// Listeners to be notified of significant happenings.
private List<KMeansListener> mListeners = new ArrayList<KMeansListener>(1);
/**
* Constructor
*
* @param coordinates two-dimensional array containing the coordinates to be clustered.
* @param k the number of desired clusters.
* @param maxIterations the maximum number of clustering iterations.
* @param randomSeed seed used with the random number generator.
*/
public BasicKMeans(double[][] coordinates, int k, int maxIterations,
long randomSeed) {
mCoordinates = coordinates;
// Can't have more clusters than coordinates.
mK = Math.min(k, mCoordinates.length);
mMaxIterations = maxIterations;
mRandomSeed = randomSeed;
}
/**
* Adds a KMeansListener to be notified of significant happenings.
*
* @param l the listener to be added.
*/
public void addKMeansListener(KMeansListener l) {
synchronized (mListeners) {
if (!mListeners.contains(l)) {
mListeners.add(l);
}
}
}
/**
* Removes a KMeansListener from the listener list.
*
* @param l the listener to be removed.
*/
public void removeKMeansListener(KMeansListener l) {
synchronized (mListeners) {
mListeners.remove(l);
}
}
/**
* Posts a message to registered KMeansListeners.
*
* @param message
*/
private void postKMeansMessage(String message) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansMessage(message);
}
}
}
}
/**
* Notifies registered listeners that k-means is complete.
*
* @param clusters the output of clustering.
* @param executionTime the number of milliseconds taken to cluster.
*/
private void postKMeansComplete(Cluster[] clusters, long executionTime) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansComplete(clusters, executionTime);
}
}
}
}
/**
* Notifies registered listeners that k-means has failed because of
* a Throwable caught in the run method.
*
* @param err
*/
private void postKMeansError(Throwable err) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansError(err);
}
}
}
}
/**
* Get the clusters computed by the algorithm. This method should
* not be called until clustering has completed successfully.
*
* @return an array of Cluster objects.
*/
public Cluster[] getClusters() {
return mClusters;
}
/**
* Run the clustering algorithm.
*/
public void run() {
try {
// Note the start time.
long startTime = System.currentTimeMillis();
postKMeansMessage("K-Means clustering started");
// Randomly initialize the cluster centers creating the
// array mProtoClusters.
initCenters();
postKMeansMessage("... centers initialized");
// Perform the initial computation of distances.
computeDistances();
// Make the initial cluster assignments.
makeAssignments();
// Number of moves in the iteration and the iteration counter.
int moves = 0, it = 0;
// Main Loop:
//
// Two stopping criteria:
// - no moves in makeAssignments
// (moves == 0)
// OR
// - the maximum number of iterations has been reached
// (it == mMaxIterations)
//
do {
// Compute the centers of the clusters that need updating.
computeCenters();
// Compute the stored distances between the updated clusters and the
// coordinates.
computeDistances();
// Make this iteration's assignments.
moves = makeAssignments();
it++;
postKMeansMessage("... iteration " + it + " moves = " + moves);
} while (moves > 0 && it < mMaxIterations);
// Transform the array of ProtoClusters to an array
// of the simpler class Cluster.
mClusters = generateFinalClusters();
long executionTime = System.currentTimeMillis() - startTime;
postKMeansComplete(mClusters, executionTime);
} catch (Throwable t) {
postKMeansError(t);
} finally {
// Clean up temporary data structures used during the algorithm.
cleanup();
}
}
/**
* Randomly select coordinates to be the initial cluster centers.
*/
private void initCenters() {
Random random = new Random(mRandomSeed);
int coordCount = mCoordinates.length;
// The array mClusterAssignments is used only to keep track of the cluster
// membership for each coordinate. The method makeAssignments() uses it
// to keep track of the number of moves.
if (mClusterAssignments == null) {
mClusterAssignments = new int[coordCount];
// Initialize to -1 to indicate that they haven't been assigned yet.
Arrays.fill(mClusterAssignments, -1);
}
// Place the coordinate indices into an array and shuffle it.
int[] indices = new int[coordCount];
for (int i = 0; i < coordCount; i++) {
indices = i;
}
for (int i = 0, m = coordCount; m > 0; i++, m--) {
int j = i + random.nextInt(m);
if (i != j) {
// Swap the indices.
indices ^= indices[j];
indices[j] ^= indices;
indices ^= indices[j];
}
}
mProtoClusters = new ProtoCluster[mK];
for (int i=0; i<mK; i++) {
int coordIndex = indices;
mProtoClusters = new ProtoCluster(mCoordinates[coordIndex], coordIndex);
mClusterAssignments[indices] = i;
}
}
/**
* Recompute the centers of the protoclusters with
* update flags set to true.
*/
private void computeCenters() {
int numClusters = mProtoClusters.length;
// Sets the update flags of the protoclusters that haven't been deleted and
// whose memberships have changed in the iteration just completed.
//
for (int c = 0; c < numClusters; c++) {
ProtoCluster cluster = mProtoClusters[c];
if (cluster.getConsiderForAssignment()) {
if (!cluster.isEmpty()) {
// This sets the protocluster's update flag to
// true only if its membership changed in last call
// to makeAssignments().
cluster.setUpdateFlag();
// If the update flag was set, update the center.
if (cluster.needsUpdate()) {
cluster.updateCenter(mCoordinates);
}
} else {
// When a cluster loses all of its members, it
// falls out of contention. So it is possible for
// k-means to return fewer than k clusters.
cluster.setConsiderForAssignment(false);
}
}
}
}
/**
* Compute distances between coodinates and cluster centers,
* storing them in the distance cache. Only distances that
* need to be computed are computed. This is determined by
* distance update flags in the protocluster objects.
*/
private void computeDistances() throws InsufficientMemoryException {
int numCoords = mCoordinates.length;
int numClusters = mProtoClusters.length;
if (mDistanceCache == null) {
// Explicit garbage collection to reduce likelihood of insufficient
// memory.
System.gc();
// Ensure there is enough memory available for the distances.
// Throw an exception if not.
long memRequired = 8L * numCoords * numClusters;
if (Runtime.getRuntime().freeMemory() < memRequired) {
throw new InsufficientMemoryException();
}
// Instantiate an array to hold the distances between coordinates
// and cluster centers
mDistanceCache = new double[numCoords][numClusters];
}
for (int coord=0; coord < numCoords; coord++) {
// Update the distances between the coordinate and all
// clusters currently in contention with update flags set.
for (int clust=0; clust<numClusters; clust++) {
ProtoCluster cluster = mProtoClusters[clust];
if (cluster.getConsiderForAssignment() && cluster.needsUpdate()) {
mDistanceCache[coord][clust] =
distance(mCoordinates[coord], cluster.getCenter());
}
}
}
}
/**
* Assign each coordinate to the nearest cluster. Called once
* per iteration. Returns the number of coordinates that have
* changed their cluster membership.
*/
private int makeAssignments() {
int moves = 0;
int coordCount = mCoordinates.length;
// Checkpoint the clusters, so we'll be able to tell
// which ones have changed after all the assignments have been
// made.
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
mProtoClusters[c].checkPoint();
}
}
// Now do the assignments.
for (int i = 0; i < coordCount; i++) {
int c = nearestCluster(i);
mProtoClusters[c].add(i);
if (mClusterAssignments != c) {
mClusterAssignments = c;
moves++;
}
}
return moves;
}
/**
* Find the nearest cluster to the coordinate identified by
* the specified index.
*/
private int nearestCluster(int ndx) {
int nearest = -1;
double min = Double.MAX_VALUE;
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
double d = mDistanceCache[ndx][c];
if (d < min) {
min = d;
nearest = c;
}
}
}
return nearest;
}
/**
* Compute the euclidean distance between the two arguments.
*/
private static double distance(double[] coord, double[] center) {
int len = coord.length;
double sumSquared = 0.0;
for (int i=0; i<len; i++) {
double v = coord - center;
sumSquared += v*v;
}
return Math.sqrt(sumSquared);
}
/**
* Generate an array of Cluster objects from mProtoClusters.
*
* @return array of Cluster object references.
*/
private Cluster[] generateFinalClusters() {
int numClusters = mProtoClusters.length;
// Convert the proto-clusters to the final Clusters.
//
// - accumulate in a list.
List<Cluster> clusterList = new ArrayList<Cluster>(numClusters);
for (int c = 0; c < numClusters; c++) {
ProtoCluster pcluster = mProtoClusters[c];
if (!pcluster.isEmpty()) {
Cluster cluster = new Cluster(pcluster.getMembership(), pcluster.getCenter());
clusterList.add(cluster);
}
}
// - convert list to an array.
Cluster[] clusters = new Cluster[clusterList.size()];
clusterList.toArray(clusters);
return clusters;
}
/**
* Clean up items used by the clustering algorithm that are no longer needed.
*/
private void cleanup() {
mProtoClusters = null;
mDistanceCache = null;
mClusterAssignments = null;
}
/**
* Cluster class used temporarily during clustering. Upon completion,
* the array of ProtoClusters is transformed into an array of
* Clusters.
*/
private static class ProtoCluster {
// The previous iteration's cluster membership and
// the current iteration's membership. Compared to see if the
// cluster has changed during the last iteration.
private int[] mPreviousMembership;
private int[] mCurrentMembership;
private int mCurrentSize;
// The cluster center.
private double[] mCenter;
// Born true, so the first call to updateDistances() will set all the
// distances.
private boolean mUpdateFlag = true;
// Whether or not this cluster takes part in the operations.
private boolean mConsiderForAssignment = true;
/**
* Constructor
*
* @param center the initial cluster center.
* @param coordIndex the initial member.
*/
ProtoCluster(double[] center, int coordIndex) {
mCenter = (double[]) center.clone();
// No previous membership.
mPreviousMembership = new int[0];
// Provide space for 10 members to be added initially.
mCurrentMembership = new int[10];
mCurrentSize = 0;
add(coordIndex);
}
/**
* Get the members of this protocluster.
*
* @return an array of coordinate indices.
*/
int[] getMembership() {
trimCurrentMembership();
return mCurrentMembership;
}
/**
* Get the protocluster's center.
*
* @return
*/
double[] getCenter() {
return mCenter;
}
/**
* Reduces the length of the array of current members to
* the number of members.
*/
void trimCurrentMembership() {
if (mCurrentMembership.length > mCurrentSize) {
int[] temp = new int[mCurrentSize];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
}
/**
* Add a coordinate to the protocluster.
*
* @param ndx index of the coordinate to be added.
*/
void add(int ndx) {
// Ensure there's space to add the new member.
if (mCurrentSize == mCurrentMembership.length) {
// If not, double the size of mCurrentMembership.
int newCapacity = Math.max(10, 2*mCurrentMembership.length);
int[] temp = new int[newCapacity];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
// Add the index.
mCurrentMembership[mCurrentSize++] = ndx;
}
/**
* Does the protocluster contain any members?
*
* @return true if the cluster is empty.
*/
boolean isEmpty() {
return mCurrentSize == 0;
}
/**
* Compares the previous and the current membership.
* Sets the update flag to true if the membership
* changed in the previous call to makeAssignments().
*/
void setUpdateFlag() {
// Trim the current membership array length down to the
// number of members.
trimCurrentMembership();
mUpdateFlag = false;
if (mPreviousMembership.length == mCurrentSize) {
for (int i=0; i<mCurrentSize; i++) {
if (mPreviousMembership != mCurrentMembership) {
mUpdateFlag = true;
break;
}
}
} else { // Number of members has changed.
mUpdateFlag = true;
}
}
/**
* Clears the current membership after copying it to the
* previous membership.
*/
void checkPoint() {
mPreviousMembership = mCurrentMembership;
mCurrentMembership = new int[10];
mCurrentSize = 0;
}
/**
* Is this protocluster currently in contention?
*
* @return true if this cluster is still in the running.
*/
boolean getConsiderForAssignment() {
return mConsiderForAssignment;
}
/**
* Set the flag to indicate that this protocluster is
* in or out of contention.
*
* @param b
*/
void setConsiderForAssignment(boolean b) {
mConsiderForAssignment = b;
}
/**
* Get the value of the update flag. This value is
* used to determine whether to update the cluster center and
* whether to recompute distances to the cluster.
*
* @return the value of the update flag.
*/
boolean needsUpdate() {
return mUpdateFlag;
}
/**
* Update the cluster center.
*
* @param coordinates the array of coordinates.
*/
void updateCenter(double[][] coordinates) {
Arrays.fill(mCenter, 0.0);
if (mCurrentSize > 0) {
for (int i=0; i<mCurrentSize; i++) {
double[] coord = coordinates[mCurrentMembership];
for (int j=0; j<coord.length; j++) {
mCenter[j] += coord[j];
}
}
for (int i=0; i<mCenter.length; i++) {
mCenter /= mCurrentSize;
}
}
}
}
}
2- BenchmarkedKMeans
package kmeans;
import java.io.*;
import java.util.*;
import javax.swing.JOptionPane;
/**
* Same as BasicKMeans, but with the addition of
* timing statements.
*/
public class BenchmarkedKMeans implements KMeans {
// Temporary clusters used during the clustering process. Converted to
// an array of the simpler class Cluster at the conclusion.
private ProtoCluster[] mProtoClusters;
// Cache of coordinate-to-cluster distances. Number of entries =
// number of clusters X number of coordinates.
private double[][] mDistanceCache;
// Used in makeAssignments() to figure out how many moves are made
// during each iteration -- the cluster assignment for coordinate n is
// found in mClusterAssignments[n] where the N coordinates are numbered
// 0 ... (N-1)
private int[] mClusterAssignments;
// 2D array holding the coordinates to be clustered.
private double[][] mCoordinates;
// The desired number of clusters and maximum number
// of iterations.
private int mK, mMaxIterations;
// Seed for the random number generator used to select
// coordinates for the initial cluster centers.
private long mRandomSeed;
// An array of Cluster objects: the output of k-means.
private Cluster[] mClusters;
// Listeners to be notified of significant happenings.
private List<KMeansListener> mListeners = new ArrayList<KMeansListener>(1);
// Fields to accumulate the time in milliseconds for
// initializing the centers, computing the distances,
// computing the centers, and making the assignments.
private long mInitCentersMS, mComputeDistancesMS,
mComputeCentersMS, mAssignmentMS;
/**
* Constructor
*
* @param coordinates two-dimensional array containing the coordinates to be clustered.
* @param k the number of desired clusters.
* @param maxIterations the maximum number of clustering iterations.
* @param randomSeed seed used with the random number generator.
*/
public BenchmarkedKMeans(double[][] coordinates, int k, int maxIterations,
long randomSeed) {
mCoordinates = coordinates;
// Can't have more clusters than coordinates.
mK = Math.min(k, mCoordinates.length);
mMaxIterations = maxIterations;
mRandomSeed = randomSeed;
}
/**
* Adds a KMeansListener to be notified of significant happenings.
*
* @param l the listener to be added.
*/
public void addKMeansListener(KMeansListener l) {
synchronized (mListeners) {
if (!mListeners.contains(l)) {
mListeners.add(l);
}
}
}
/**
* Removes a KMeansListener from the listener list.
*
* @param l the listener to be removed.
*/
public void removeKMeansListener(KMeansListener l) {
synchronized (mListeners) {
mListeners.remove(l);
}
}
/**
* Posts a message to registered KMeansListeners.
*
* @param message
*/
private void postKMeansMessage(String message) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansMessage(message);
}
}
}
}
/**
* Notifies registered listeners that k-means is complete.
*
* @param clusters the output of clustering.
* @param executionTime the number of milliseconds taken to cluster.
*/
private void postKMeansComplete(Cluster[] clusters, long executionTime) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansComplete(clusters, executionTime);
}
}
}
}
/**
* Notifies registered listeners that k-means has failed because of
* a Throwable caught in the run method.
*
* @param err
*/
private void postKMeansError(Throwable err) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansError(err);
}
}
}
}
/**
* Get the clusters computed by the algorithm. This method should
* not be called until clustering has completed successfully.
*
* @return an array of Cluster objects.
*/
public Cluster[] getClusters() {
return mClusters;
}
/**
* Run the clustering algorithm.
*/
public void run() {
try {
// Note the start time.
long startTime = System.currentTimeMillis();
postKMeansMessage("K-Means clustering started");
// Randomly initialize the cluster centers creating the
// array mProtoClusters.
initCenters();
postKMeansMessage("... centers initialized");
// Perform the initial computation of distances.
computeDistances();
// Make the initial cluster assignments.
makeAssignments();
// Number of moves in the iteration and the iteration counter.
int moves = 0, it = 0;
// Main Loop:
//
// Two stopping criteria:
// - no moves in makeAssignments
// (moves == 0)
// OR
// - the maximum number of iterations has been reached
// (it == mMaxIterations)
//
do {
// Compute the centers of the clusters that need updating.
computeCenters();
// Compute the stored distances between the updated clusters and the
// coordinates.
computeDistances();
// Make this iteration's assignments.
moves = makeAssignments();
it++;
postKMeansMessage("... iteration " + it + " moves = " + moves);
} while (moves > 0 && it < mMaxIterations);
// Transform the array of ProtoClusters to an array
// of the simpler class Cluster.
mClusters = generateFinalClusters();
long executionTime = System.currentTimeMillis() - startTime;
// Post the time statistics.
postKMeansMessage("TIME STATISTICS:\nInitializing centers: " +
percentString(mInitCentersMS, executionTime));
postKMeansMessage(" Computing centers: " +
percentString(mComputeCentersMS, executionTime));
postKMeansMessage(" Computing distances: " +
percentString(mComputeDistancesMS, executionTime));
postKMeansMessage(" Making assignments: " +
percentString(mAssignmentMS, executionTime));
postKMeansComplete(mClusters, executionTime);
} catch (Throwable t) {
postKMeansError(t);
} finally {
// Clean up temporary data structures used during the algorithm.
cleanup();
}
}
/**
* Convenience method for generating percentage string.
*
* @param numerator
* @param denominator
*
* @return a string of the form 10.2%
*/
private static String percentString(long numerator, long denominator) {
StringWriter sw = new StringWriter();
PrintWriter pw = new PrintWriter(sw);
double percent = (100.0*numerator)/denominator;
pw.printf("%02.1f", percent);
return sw.toString();
}
/**
* Randomly select coordinates to be the initial cluster centers.
*/
private void initCenters() {
long t = System.currentTimeMillis();
Random random = new Random(mRandomSeed);
int coordCount = mCoordinates.length;
// The array mClusterAssignments is used only to keep track of the cluster
// membership for each coordinate. The method makeAssignments() uses it
// to keep track of the number of moves.
if (mClusterAssignments == null) {
mClusterAssignments = new int[coordCount];
// Initialize to -1 to indicate that they haven't been assigned yet.
Arrays.fill(mClusterAssignments, -1);
}
// Place the coordinate indices into an array and shuffle it.
int[] indices = new int[coordCount];
for (int i = 0; i < coordCount; i++) {
indices = i;
}
for (int i = 0, m = coordCount; m > 0; i++, m--) {
int j = i + random.nextInt(m);
if (i != j) {
// Swap the indices.
indices ^= indices[j];
indices[j] ^= indices;
indices ^= indices[j];
}
}
mProtoClusters = new ProtoCluster[mK];
for (int i=0; i<mK; i++) {
int coordIndex = indices;
mProtoClusters = new ProtoCluster(mCoordinates[coordIndex], coordIndex);
mClusterAssignments[indices] = i;
}
mInitCentersMS += (System.currentTimeMillis() - t);
}
/**
* Recompute the centers of the protoclusters with
* update flags set to true.
*/
private void computeCenters() {
long t = System.currentTimeMillis();
int numClusters = mProtoClusters.length;
// Sets the update flags of the protoclusters that haven't been deleted and
// whose memberships have changed in the iteration just completed.
//
for (int c = 0; c < numClusters; c++) {
ProtoCluster cluster = mProtoClusters[c];
if (cluster.getConsiderForAssignment()) {
if (!cluster.isEmpty()) {
// This sets the protocluster's update flag to
// true only if its membership changed in last call
// to makeAssignments().
cluster.setUpdateFlag();
// If the update flag was set, update the center.
if (cluster.needsUpdate()) {
cluster.updateCenter(mCoordinates);
}
} else {
// When a cluster loses all of its members, it
// falls out of contention. So it is possible for
// k-means to return fewer than k clusters.
cluster.setConsiderForAssignment(false);
}
}
}
mComputeCentersMS += (System.currentTimeMillis() - t);
}
/**
* Compute distances between coodinates and cluster centers,
* storing them in the distance cache. Only distances that
* need to be computed are computed. This is determined by
* distance update flags in the protocluster objects.
*/
private void computeDistances() throws InsufficientMemoryException {
long t = System.currentTimeMillis();
int numCoords = mCoordinates.length;
int numClusters = mProtoClusters.length;
if (mDistanceCache == null) {
// Explicit garbage collection to reduce likelihood of insufficient
// memory.
System.gc();
// Ensure there is enough memory available for the distances.
// Throw an exception if not.
long memRequired = 8L * numCoords * numClusters;
if (Runtime.getRuntime().freeMemory() < memRequired) {
throw new InsufficientMemoryException();
}
// Instantiate an array to hold the distances between coordinates
// and cluster centers
mDistanceCache = new double[numCoords][numClusters];
}
for (int coord=0; coord < numCoords; coord++) {
// Update the distances between the coordinate and all
// clusters currently in contention with update flags set.
for (int clust=0; clust<numClusters; clust++) {
ProtoCluster cluster = mProtoClusters[clust];
if (cluster.getConsiderForAssignment() && cluster.needsUpdate()) {
mDistanceCache[coord][clust] =
distance(mCoordinates[coord], cluster.getCenter());
}
}
}
mComputeDistancesMS += (System.currentTimeMillis() - t);
}
/**
* Assign each coordinate to the nearest cluster. Called once
* per iteration. Returns the number of coordinates that have
* changed their cluster membership.
*/
private int makeAssignments() {
long t = System.currentTimeMillis();
int moves = 0;
int coordCount = mCoordinates.length;
// Checkpoint the clusters, so we'll be able to tell
// which ones have changed after all the assignments have been
// made.
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
mProtoClusters[c].checkPoint();
}
}
// Now do the assignments.
for (int i = 0; i < coordCount; i++) {
int c = nearestCluster(i);
mProtoClusters[c].add(i);
if (mClusterAssignments != c) {
mClusterAssignments = c;
moves++;
}
}
mAssignmentMS += (System.currentTimeMillis() - t);
return moves;
}
/**
* Find the nearest cluster to the coordinate identified by
* the specified index.
*/
private int nearestCluster(int ndx) {
int nearest = -1;
double min = Double.MAX_VALUE;
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
double d = mDistanceCache[ndx][c];
if (d < min) {
min = d;
nearest = c;
}
}
}
return nearest;
}
/**
* Compute the euclidean distance between the two arguments.
*/
private static double distance(double[] coord, double[] center) {
int len = coord.length;
double sumSquared = 0.0;
for (int i=0; i<len; i++) {
double v = coord - center;
sumSquared += v*v;
}
return Math.sqrt(sumSquared);
}
/**
* Generate an array of Cluster objects from mProtoClusters.
*
* @return array of Cluster object references.
*/
private Cluster[] generateFinalClusters() {
int numClusters = mProtoClusters.length;
// Convert the proto-clusters to the final Clusters.
//
// - accumulate in a list.
List<Cluster> clusterList = new ArrayList<Cluster>(numClusters);
for (int c = 0; c < numClusters; c++) {
ProtoCluster pcluster = mProtoClusters[c];
if (!pcluster.isEmpty()) {
Cluster cluster = new Cluster(pcluster.getMembership(), pcluster.getCenter());
clusterList.add(cluster);
}
}
// - convert list to an array.
Cluster[] clusters = new Cluster[clusterList.size()];
clusterList.toArray(clusters);
return clusters;
}
/**
* Clean up items used by the clustering algorithm that are no longer needed.
*/
private void cleanup() {
mProtoClusters = null;
mDistanceCache = null;
mClusterAssignments = null;
}
/**
* Cluster class used temporarily during clustering. Upon completion,
* the array of ProtoClusters is transformed into an array of
* Clusters.
*/
private static class ProtoCluster {
// The previous iteration's cluster membership and
// the current iteration's membership. Compared to see if the
// cluster has changed during the last iteration.
private int[] mPreviousMembership;
private int[] mCurrentMembership;
private int mCurrentSize;
// The cluster center.
private double[] mCenter;
// Born true, so the first call to updateDistances() will set all the
// distances.
private boolean mUpdateFlag = true;
// Whether or not this cluster takes part in the operations.
private boolean mConsiderForAssignment = true;
/**
* Constructor
*
* @param center the initial cluster center.
* @param coordIndex the initial member.
*/
ProtoCluster(double[] center, int coordIndex) {
mCenter = (double[]) center.clone();
// No previous membership.
mPreviousMembership = new int[0];
// Provide space for 10 members to be added initially.
mCurrentMembership = new int[10];
mCurrentSize = 0;
add(coordIndex);
}
/**
* Get the members of this protocluster.
*
* @return an array of coordinate indices.
*/
int[] getMembership() {
trimCurrentMembership();
return mCurrentMembership;
}
/**
* Get the protocluster's center.
*
* @return
*/
double[] getCenter() {
return mCenter;
}
/**
* Reduces the length of the array of current members to
* the number of members.
*/
void trimCurrentMembership() {
if (mCurrentMembership.length > mCurrentSize) {
int[] temp = new int[mCurrentSize];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
}
/**
* Add a coordinate to the protocluster.
*
* @param ndx index of the coordinate to be added.
*/
void add(int ndx) {
// Ensure there's space to add the new member.
if (mCurrentSize == mCurrentMembership.length) {
// If not, double the size of mCurrentMembership.
int newCapacity = Math.max(10, 2*mCurrentMembership.length);
int[] temp = new int[newCapacity];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
// Add the index.
mCurrentMembership[mCurrentSize++] = ndx;
}
/**
* Does the protocluster contain any members?
*
* @return true if the cluster is empty.
*/
boolean isEmpty() {
return mCurrentSize == 0;
}
/**
* Compares the previous and the current membership.
* Sets the update flag to true if the membership
* changed in the previous call to makeAssignments().
*/
void setUpdateFlag() {
// Trim the current membership array length down to the
// number of members.
trimCurrentMembership();
mUpdateFlag = false;
if (mPreviousMembership.length == mCurrentSize) {
for (int i=0; i<mCurrentSize; i++) {
if (mPreviousMembership != mCurrentMembership) {
mUpdateFlag = true;
break;
}
}
} else { // Number of members has changed.
mUpdateFlag = true;
}
}
/**
* Clears the current membership after copying it to the
* previous membership.
*/
void checkPoint() {
mPreviousMembership = mCurrentMembership;
mCurrentMembership = new int[10];
mCurrentSize = 0;
}
/**
* Is this protocluster currently in contention?
*
* @return true if this cluster is still in the running.
*/
boolean getConsiderForAssignment() {
return mConsiderForAssignment;
}
/**
* Set the flag to indicate that this protocluster is
* in or out of contention.
*
* @param b
*/
void setConsiderForAssignment(boolean b) {
mConsiderForAssignment = b;
}
/**
* Get the value of the update flag. This value is
* used to determine whether to update the cluster center and
* whether to recompute distances to the cluster.
*
* @return the value of the update flag.
*/
boolean needsUpdate() {
return mUpdateFlag;
}
/**
* Update the cluster center.
*
* @param coordinates the array of coordinates.
*/
void updateCenter(double[][] coordinates) {
Arrays.fill(mCenter, 0.0);
if (mCurrentSize > 0) {
for (int i=0; i<mCurrentSize; i++) {
double[] coord = coordinates[mCurrentMembership];
for (int j=0; j<coord.length; j++) {
mCenter[j] += coord[j];
}
}
for (int i=0; i<mCenter.length; i++) {
mCenter /= mCurrentSize;
}
}
}
}
}
3- Cluster
package kmeans;
/**
* Class to represent a cluster of coordinates.
*/
public class Cluster {
// Indices of the member coordinates.
private int[] mMemberIndexes;
// The cluster center.
private double[] mCenter;
/**
* Constructor.
*
* @param memberIndexes indices of the member coordinates.
* @param center the cluster center.
*/
public Cluster(int[] memberIndexes, double[] center) {
mMemberIndexes = memberIndexes;
mCenter = center;
}
/**
* Get the member indices.
*
* @return an array containing the indices of the member coordinates.
*/
public int[] getMemberIndexes()
{
return mMemberIndexes;
}
/**
* Get the cluster center.
*
* @return a reference to the cluster center array.
*/
public double[] getCenter()
{
return mCenter;
}
}
4- ConcurrentKMeans
package kmeans;
import java.util.*;
import java.util.concurrent.*;
/**
* The version of K-means clustering adapted for true concurrency
* or simultaneous multithreading (SMT). The subtasks of
* computing distances and making assignments are delegate to
* a subtask manager which oversees a thread pool.
*/
public class ConcurrentKMeans implements KMeans {
// Temporary clusters used during the clustering process. Converted to
// an array of the simpler class Cluster at the conclusion.
private ProtoCluster[] mProtoClusters;
// Cache of coordinate-to-cluster distances. Number of entries =
// number of clusters X number of coordinates.
private double[][] mDistanceCache;
// Used in makeAssignments() to figure out how many moves are made
// during each iteration -- the cluster assignment for coordinate n is
// found in mClusterAssignments[n] where the N coordinates are numbered
// 0 ... (N-1)
private int[] mClusterAssignments;
// 2D array holding the coordinates to be clustered.
private double[][] mCoordinates;
// The desired number of clusters and maximum number
// of iterations.
private int mK, mMaxIterations;
// Seed for the random number generator used to select
// coordinates for the initial cluster centers.
private long mRandomSeed;
// The number of threads used to perform the subtasks.
private int mThreadCount;
// Subtask manager that handles the thread pool to which
// time-consuming tasks are delegated.
private SubtaskManager mSubtaskManager;
// An array of Cluster objects: the output of k-means.
private Cluster[] mClusters;
// Listeners to be notified of significant happenings.
private List<KMeansListener> mListeners = new ArrayList<KMeansListener>(1);
/**
* Constructor
*
* @param coordinates two-dimensional array containing the coordinates to be clustered.
* @param k the number of desired clusters.
* @param maxIterations the maximum number of clustering iterations.
* @param randomSeed seed used with the random number generator.
* @param threadCount the number of threads to be used for computing time-consuming steps.
*/
public ConcurrentKMeans(double[][] coordinates, int k, int maxIterations,
long randomSeed, int threadCount) {
mCoordinates = coordinates;
// Can't have more clusters than coordinates.
mK = Math.min(k, mCoordinates.length);
mMaxIterations = maxIterations;
mRandomSeed = randomSeed;
mThreadCount = threadCount;
}
/**
* Constructor that uses the return from
* <tt>Runtime.getRuntime().availableProcessors()</tt> as the number
* of threads for time-consuming steps.
*
* @param coordinates two-dimensional array containing the coordinates to be clustered.
* @param k the number of desired clusters.
* @param maxIterations the maximum number of clustering iterations.
* @param randomSeed seed used with the random number generator.
*/
public ConcurrentKMeans(double[][] coordinates, int k, int maxIterations,
long randomSeed) {
this (coordinates, k, maxIterations, randomSeed,
Runtime.getRuntime().availableProcessors());
}
/**
* Adds a KMeansListener to be notified of significant happenings.
*
* @param l the listener to be added.
*/
public void addKMeansListener(KMeansListener l) {
synchronized (mListeners) {
if (!mListeners.contains(l)) {
mListeners.add(l);
}
}
}
/**
* Removes a KMeansListener
*
* @param l the listener to be removed.
*/
public void removeKMeansListener(KMeansListener l) {
synchronized (mListeners) {
mListeners.remove(l);
}
}
/**
* Posts a message to registered KMeansListeners.
*
* @param message
*/
private void postKMeansMessage(String message) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansMessage(message);
}
}
}
}
/**
* Notifies registered listeners that k-means is complete.
*
* @param clusters the output of clustering.
* @param executionTime the number of milliseconds taken to cluster.
*/
private void postKMeansComplete(Cluster[] clusters, long executionTime) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansComplete(clusters, executionTime);
}
}
}
}
/**
* Notifies registered listeners that k-means has failed because of
* a Throwable caught in the run method.
*
* @param err
*/
private void postKMeansError(Throwable err) {
if (mListeners.size() > 0) {
synchronized (mListeners) {
int sz = mListeners.size();
for (int i=0; i<sz; i++) {
mListeners.get(i).kmeansError(err);
}
}
}
}
/**
* Get the clusters computed by the algorithm. This method should
* not be called until clustering has completed successfully.
*
* @return an array of Cluster objects.
*/
public Cluster[] getClusters() {
return mClusters;
}
/**
* Run the clustering algorithm.
*/
public void run() {
try {
// Note the start time.
long startTime = System.currentTimeMillis();
postKMeansMessage("K-Means clustering started");
// Randomly initialize the cluster centers creating the
// array mProtoClusters.
initCenters();
postKMeansMessage("... centers initialized");
// Instantiate the subtask manager.
mSubtaskManager = new SubtaskManager(mThreadCount);
// Post a message about the state of concurrent subprocessing.
if (mThreadCount > 1) {
postKMeansMessage("... concurrent processing mode with "
+ mThreadCount + " subtask threads");
} else {
postKMeansMessage("... non-concurrent processing mode");
}
// Perform the initial computation of distances.
computeDistances();
// Make the initial cluster assignments.
makeAssignments();
// Number of moves in the iteration and the iteration counter.
int moves = 0, it = 0;
// Main Loop:
//
// Two stopping criteria:
// - no moves in makeAssignments
// (moves == 0)
// OR
// - the maximum number of iterations has been reached
// (it == mMaxIterations)
//
do {
// Compute the centers of the clusters that need updating.
computeCenters();
// Compute the stored distances between the updated clusters and the
// coordinates.
computeDistances();
// Make this iteration's assignments.
moves = makeAssignments();
it++;
postKMeansMessage("... iteration " + it + " moves = " + moves);
} while (moves > 0 && it < mMaxIterations);
// Transform the array of ProtoClusters to an array
// of the simpler class Cluster.
mClusters = generateFinalClusters();
long executionTime = System.currentTimeMillis() - startTime;
postKMeansComplete(mClusters, executionTime);
} catch (Throwable t) {
postKMeansError(t);
} finally {
// Clean up temporary data structures used during the algorithm.
cleanup();
}
}
/**
* Randomly select coordinates to be the initial cluster centers.
*/
private void initCenters() {
Random random = new Random(mRandomSeed);
int coordCount = mCoordinates.length;
// The array mClusterAssignments is used only to keep track of the cluster
// membership for each coordinate. The method makeAssignments() uses it
// to keep track of the number of moves.
if (mClusterAssignments == null) {
mClusterAssignments = new int[coordCount];
// Initialize to -1 to indicate that they haven't been assigned yet.
Arrays.fill(mClusterAssignments, -1);
}
// Place the coordinate indices into an array and shuffle it.
int[] indices = new int[coordCount];
for (int i = 0; i < coordCount; i++) {
indices = i;
}
for (int i = 0, m = coordCount; m > 0; i++, m--) {
int j = i + random.nextInt(m);
if (i != j) {
// Swap the indices.
indices ^= indices[j];
indices[j] ^= indices;
indices ^= indices[j];
}
}
mProtoClusters = new ProtoCluster[mK];
for (int i=0; i<mK; i++) {
int coordIndex = indices;
mProtoClusters = new ProtoCluster(mCoordinates[coordIndex], coordIndex);
mClusterAssignments[indices] = i;
}
}
/**
* Recompute the centers of the protoclusters with
* update flags set to true.
*/
private void computeCenters() {
int numClusters = mProtoClusters.length;
// Sets the update flags of the protoclusters that haven't been deleted and
// whose memberships have changed in the iteration just completed.
//
for (int c = 0; c < numClusters; c++) {
ProtoCluster cluster = mProtoClusters[c];
if (cluster.getConsiderForAssignment()) {
if (!cluster.isEmpty()) {
// This sets the protocluster's update flag to
// true only if its membership changed in last call
// to makeAssignments().
cluster.setUpdateFlag();
// If the update flag was set, update the center.
if (cluster.needsUpdate()) {
cluster.updateCenter(mCoordinates);
}
} else {
// When a cluster loses all of its members, it
// falls out of contention. So it is possible for
// k-means to return fewer than k clusters.
cluster.setConsiderForAssignment(false);
}
}
}
}
/**
* Compute distances between coodinates and cluster centers,
* storing them in the distance cache. Only distances that
* need to be computed are computed. This is determined by
* distance update flags in the protocluster objects.
*/
private void computeDistances() throws InsufficientMemoryException {
if (mDistanceCache == null) {
int numCoords = mCoordinates.length;
int numClusters = mProtoClusters.length;
// Explicit garbage collection to reduce likelihood of insufficient
// memory.
System.gc();
// Ensure there is enough memory available for the distances.
// Throw an exception if not.
long memRequired = 8L * numCoords * numClusters;
if (Runtime.getRuntime().freeMemory() < memRequired) {
throw new InsufficientMemoryException();
}
// Instantiate an array to hold the distances between coordinates
// and cluster centers
mDistanceCache = new double[numCoords][numClusters];
}
// Bulk of the work is delegated to the
// SubtaskManager.
mSubtaskManager.computeDistances();
}
/**
* Assign each coordinate to the nearest cluster. Called once
* per iteration. Returns the number of coordinates that have
* changed their cluster membership.
*/
private int makeAssignments() {
// Checkpoint the clusters, so we'll be able to tell
// which one have changed after all the assignments have been
// made.
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
mProtoClusters[c].checkPoint();
}
}
// Bulk of the work is delegated to the SubtaskManager.
mSubtaskManager.makeAssignments();
// Get the number of moves from the SubtaskManager.
return mSubtaskManager.numberOfMoves();
}
/**
* Find the nearest cluster to the coordinate identified by
* the specified index.
*/
private int nearestCluster(int ndx) {
int nearest = -1;
double min = Double.MAX_VALUE;
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
if (mProtoClusters[c].getConsiderForAssignment()) {
double d = mDistanceCache[ndx][c];
if (d < min) {
min = d;
nearest = c;
}
}
}
return nearest;
}
/**
* Compute the euclidean distance between the two arguments.
*/
private static double distance(double[] coord, double[] center) {
int len = coord.length;
double sumSquared = 0.0;
for (int i=0; i<len; i++) {
double v = coord - center;
sumSquared += v*v;
}
return Math.sqrt(sumSquared);
}
/**
* Generate an array of Cluster objects from mProtoClusters.
*
* @return array of Cluster object references.
*/
private Cluster[] generateFinalClusters() {
int numClusters = mProtoClusters.length;
// Convert the proto-clusters to the final Clusters.
//
// - accumulate in a list.
List<Cluster> clusterList = new ArrayList<Cluster>(numClusters);
for (int c = 0; c < numClusters; c++) {
ProtoCluster pcluster = mProtoClusters[c];
if (!pcluster.isEmpty()) {
Cluster cluster = new Cluster(pcluster.getMembership(), pcluster.getCenter());
clusterList.add(cluster);
}
}
// - convert list to an array.
Cluster[] clusters = new Cluster[clusterList.size()];
clusterList.toArray(clusters);
return clusters;
}
/**
* Clean up items used by the clustering algorithm that are no longer needed.
*/
private void cleanup() {
mProtoClusters = null;
mDistanceCache = null;
mClusterAssignments = null;
if (mSubtaskManager != null) {
mSubtaskManager.shutdown();
mSubtaskManager = null;
}
}
/**
* Cluster class used temporarily during clustering. Upon completion,
* the array of ProtoClusters is transformed into an array of
* Clusters.
*/
private static class ProtoCluster {
// The previous iteration's cluster membership and
// the current iteration's membership. Compared to see if the
// cluster has changed during the last iteration.
private int[] mPreviousMembership;
private int[] mCurrentMembership;
private int mCurrentSize;
// The cluster center.
private double[] mCenter;
// Born true, so the first call to updateDistances() will set all the
// distances.
private boolean mUpdateFlag = true;
// Whether or not this cluster takes part in the operations.
private boolean mConsiderForAssignment = true;
/**
* Constructor
*
* @param center the initial cluster center.
* @param coordIndex the initial member.
*/
ProtoCluster(double[] center, int coordIndex) {
mCenter = (double[]) center.clone();
// No previous membership.
mPreviousMembership = new int[0];
// Provide space for 10 members to be added initially.
mCurrentMembership = new int[10];
mCurrentSize = 0;
add(coordIndex);
}
/**
* Get the members of this protocluster.
*
* @return an array of coordinate indices.
*/
int[] getMembership() {
trimCurrentMembership();
return mCurrentMembership;
}
/**
* Get the protocluster's center.
*
* @return
*/
double[] getCenter() {
return mCenter;
}
/**
* Reduces the length of the array of current members to
* the number of members.
*/
void trimCurrentMembership() {
if (mCurrentMembership.length > mCurrentSize) {
int[] temp = new int[mCurrentSize];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
}
/**
* Add a coordinate to the protocluster. Note that this
* method has to be synchronized, because multiple threads
* may be adding members to the cluster.
*
* @param ndx index of the coordinate to be added.
*/
synchronized void add(int ndx) {
// Ensure there's space to add the new member.
if (mCurrentSize == mCurrentMembership.length) {
// If not, double the size of mCurrentMembership.
int newCapacity = Math.max(10, 2*mCurrentMembership.length);
int[] temp = new int[newCapacity];
System.arraycopy(mCurrentMembership, 0, temp, 0, mCurrentSize);
mCurrentMembership = temp;
}
// Add the index.
mCurrentMembership[mCurrentSize++] = ndx;
}
/**
* Does the protocluster contain any members?
*
* @return true if the cluster is empty.
*/
boolean isEmpty() {
return mCurrentSize == 0;
}
/**
* Compares the previous and the current membership.
* Sets the update flag to true if the membership
* changed in the previous call to makeAssignments().
*/
void setUpdateFlag() {
// Trim the current membership array length down to the
// number of members.
trimCurrentMembership();
// Since members may have been added by multiple threads, they
// are probably not in order. They must be sorted in order to
// do a valid comparison with mPreviousMembership.
Arrays.sort(mCurrentMembership);
mUpdateFlag = false;
if (mPreviousMembership.length == mCurrentSize) {
for (int i=0; i<mCurrentSize; i++) {
if (mPreviousMembership != mCurrentMembership) {
mUpdateFlag = true;
break;
}
}
} else { // Number of members has changed.
mUpdateFlag = true;
}
}
/**
* Clears the current membership after copying it to the
* previous membership.
*/
void checkPoint() {
mPreviousMembership = mCurrentMembership;
mCurrentMembership = new int[10];
mCurrentSize = 0;
}
/**
* Is this protocluster currently in contention?
*
* @return true if this cluster is still in the running.
*/
boolean getConsiderForAssignment() {
return mConsiderForAssignment;
}
/**
* Set the flag to indicate that this protocluster is
* in or out of contention.
*
* @param b
*/
void setConsiderForAssignment(boolean b) {
mConsiderForAssignment = b;
}
/**
* Get the value of the update flag. This value is
* used to determine whether to update the cluster center and
* whether to recompute distances to the cluster.
*
* @return the value of the update flag.
*/
boolean needsUpdate() {
return mUpdateFlag;
}
/**
* Update the cluster center.
*
* @param coordinates the array of coordinates.
*/
void updateCenter(double[][] coordinates) {
Arrays.fill(mCenter, 0.0);
if (mCurrentSize > 0) {
for (int i=0; i<mCurrentSize; i++) {
double[] coord = coordinates[mCurrentMembership];
for (int j=0; j<coord.length; j++) {
mCenter[j] += coord[j];
}
}
for (int i=0; i<mCenter.length; i++) {
mCenter /= mCurrentSize;
}
}
}
}
/**
* The class which manages the SMT-adapted subtasks.
*/
private class SubtaskManager {
// Codes used to identify what step is being done.
static final int DOING_NOTHING = 0;
static final int COMPUTING_DISTANCES = 1;
static final int MAKING_ASSIGNMENTS = 2;
// What the object is currently doing. Set to one of the
// three codes above.
private int mDoing = DOING_NOTHING;
// True if the at least one of the Workers is doing something.
private boolean mWorking;
// The executor that runs the Workers.
// When in multiple processor mode, this is a ThreadPoolExecutor
// with a fixed number of threads. In single-processor mode, it's
// a simple implementation that calls the single worker's run
// method directly.
private Executor mExecutor;
// A Barrier to wait on multiple Workers to finish up the current task.
// In single-processor mode, there is no need for a barrier, so it
// is not set.
private CyclicBarrier mBarrier;
// The worker objects which implement Runnable.
private Worker[] mWorkers;
/**
* Constructor
*
* @param numThreads the number of worker threads to be used for
* the subtasks.
*/
SubtaskManager(int numThreads) {
if (numThreads <= 0) {
throw new IllegalArgumentException("number of threads <= 0: "
+ numThreads);
}
int coordCount = mCoordinates.length;
// There would be no point in having more workers than
// coordinates, since some of the workers would have nothing
// to do.
if (numThreads > coordCount) {
numThreads = coordCount;
}
// Create the workers.
mWorkers = new Worker[numThreads];
// To hold the number of coordinates for each worker.
int[] coordsPerWorker = new int[numThreads];
// Initialize with the base amounts.
Arrays.fill(coordsPerWorker, coordCount/numThreads);
// There may be some leftovers, since coordCount may not be
// evenly divisible by numWorkers. Add a coordinate to each
// until all are covered.
int leftOvers = coordCount - numThreads * coordsPerWorker[0];
for (int i = 0; i < leftOvers; i++) {
coordsPerWorker++;
}
int startCoord = 0;
// Instantiate the workers.
for (int i = 0; i < numThreads; i++) {
// Each worker needs to know its starting coordinate and the number of
// coordinates it handles.
mWorkers = new Worker(startCoord, coordsPerWorker);
startCoord += coordsPerWorker;
}
if (numThreads == 1) { // Single-processor mode.
// Create a simple executor that directly calls the single
// worker's run method. Do not set the barrier.
mExecutor = new Executor() {
public void execute(Runnable runnable) {
if (!Thread.interrupted()) {
runnable.run();
} else {
throw new RejectedExecutionException();
}
}
};
} else { // Multiple-processor mode.
// Need the barrier to notify the controlling thread when the
// Workers are done.
mBarrier = new CyclicBarrier(numThreads, new Runnable() {
public void run() {
// Method called after all workers haved called await() on the
// barrier. The call to workersDone()
// unblocks the controlling thread.
workersDone();
}
});
// Set the executor to a fixed thread pool with
// threads that do not time out.
mExecutor = Executors.newFixedThreadPool(numThreads);
}
}
/**
* Make the cluster assignments.
*
* @return true if nothing went wrong.
*/
boolean makeAssignments() {
mDoing = MAKING_ASSIGNMENTS;
return work();
}
/**
* Compute the distances between the coordinates and those centers with
* update flags.
*
* @return true if nothing went wrong.
*/
boolean computeDistances() {
mDoing = COMPUTING_DISTANCES;
return work();
}
/**
* Perform the current subtask, waiting until all the workers
* finish their part of the current task before returning.
*
* @return true if the subtask succeeded.
*/
private boolean work() {
boolean ok = false;
// Set the working flag to true.
mWorking = true;
try {
if (mBarrier != null) {
// Resets the barrier so it can be reused if
// this is not the first call to this method.
mBarrier.reset();
}
// Now execute the run methods on the Workers.
for (int i = 0; i < mWorkers.length; i++) {
mExecutor.execute(mWorkers);
}
if (mBarrier != null) {
// Block until the workers are done. The barrier
// triggers the unblocking.
waitOnWorkers();
// If the isBroken() method of the barrier returns false,
// no problems.
ok = !mBarrier.isBroken();
} else {
// No barrier, so the run() method of a single worker
// was called directly and everything must have worked
// if we made it here.
ok = true;
}
} catch (RejectedExecutionException ree) {
// Possibly thrown by the executor.
} finally {
mWorking = false;
}
return ok;
}
/**
* Called from work() to put the controlling thread into
* wait mode until the barrier calls workersDone().
*/
private synchronized void waitOnWorkers() {
// It is possible for the workers to have finished so quickly that
// workersDone() has already been called. Since workersDone() sets
// mWorking to false, check this flag before going into wait mode.
// Not doing so could result in hanging the SubtaskManager.
while (mWorking) {
try {
// Blocks until workersDone() is called.
wait();
} catch (InterruptedException ie) {
// mBarrier.isBroken() will return true.
break;
}
}
}
/**
* Notifies the controlling thread that it can come out of
* wait mode.
*/
private synchronized void workersDone() {
// If this gets called before waitOnWorkers(), setting this
// to false prevents waitOnWorkers() from entering a
// permanent wait.
mWorking = false;
notifyAll();
}
/**
* Shutdown the thread pool when k-means is finished.
*/
void shutdown() {
if (mExecutor instanceof ThreadPoolExecutor) {
// This terminates the threads in the thread pool.
((ThreadPoolExecutor) mExecutor).shutdownNow();
}
}
/**
* Returns the number of cluster assignment changes made in the
* previous call to makeAssignments().
*/
int numberOfMoves() {
// Sum the contributions from the workers.
int moves = 0;
for (int i=0; i<mWorkers.length; i++) {
moves += mWorkers.numberOfMoves();
}
return moves;
}
/**
* The class which does the hard work of the subtasks.
*/
private class Worker implements Runnable {
// Defines range of coordinates to cover.
private int mStartCoord, mNumCoords;
// Number of moves made by this worker in the last call
// to workerMakeAssignments(). The SubtaskManager totals up
// this value from all the workers in numberOfMoves().
private int mMoves;
/**
* Constructor
*
* @param startCoord index of the first coordinate covered by
* this Worker.
* @param numCoords the number of coordinates covered.
*/
Worker(int startCoord, int numCoords) {
mStartCoord = startCoord;
mNumCoords = numCoords;
}
/**
* Returns the number of moves this worker made in the last
* execution of workerMakeAssignments()
*/
int numberOfMoves() {
return mMoves;
}
/**
* The run method. It accesses the SubtaskManager field mDoing
* to determine what subtask to perform.
*/
public void run() {
try {
switch (mDoing) {
case COMPUTING_DISTANCES:
workerComputeDistances();
break;
case MAKING_ASSIGNMENTS:
workerMakeAssignments();
break;
}
} finally {
// If there's a barrier, call its await() method. To ensure it
// gets done, it's placed in the finally clause.
if (mBarrier != null) {
try {
mBarrier.await();
// barrier.isBroken() will return true if either of these
// exceptions happens, so the SubtaskManager will detect
// the problem.
} catch (InterruptedException ex) {
} catch (BrokenBarrierException ex) {
}
}
}
}
/**
* Compute the distances for the covered coordinates
* to the updated centers.
*/
private void workerComputeDistances() {
int lim = mStartCoord + mNumCoords;
for (int i = mStartCoord; i < lim; i++) {
int numClusters = mProtoClusters.length;
for (int c = 0; c < numClusters; c++) {
ProtoCluster cluster = mProtoClusters[c];
if (cluster.getConsiderForAssignment() && cluster.needsUpdate()) {
mDistanceCache[c] = distance(mCoordinates, cluster.getCenter());
}
}
}
}
/**
* Assign each covered coordinate to the nearest cluster.
*/
private void workerMakeAssignments() {
mMoves = 0;
int lim = mStartCoord + mNumCoords;
for (int i = mStartCoord; i < lim; i++) {
int c = nearestCluster(i);
mProtoClusters[c].add(i);
if (mClusterAssignments != c) {
mClusterAssignments = c;
mMoves++;
}
}
}
}
}
}
5- InsufficientMemoryException
package kmeans;
/**
* Exception thrown when insufficient memory is available to
* perform an operation. Designed to be throw before doing
* something that would cause a <code>java.lang.OutOfMemoryError</code>.
*/
public class InsufficientMemoryException extends Exception {
/**
* Constructor.
*
* @param message an explanatory message.
*/
public InsufficientMemoryException(String message) {
super(message);
}
/**
* Default constructor.
*/
public InsufficientMemoryException() {}
}
6- KMeans
package kmeans;
/**
* Simple K-Means clustering interface.
*/
public interface KMeans extends Runnable {
/**
* Adds a KMeansListener to be notified of significant happenings.
*
* @param l the listener to be added.
*/
public void addKMeansListener(KMeansListener l);
/**
* Removes a KMeansListener from the listener list.
*
* @param l the listener to be removed.
*/
public void removeKMeansListener(KMeansListener l);
/**
* Get the clusters computed by the algorithm. This method should
* not be called until clustering has completed successfully.
*
* @return an array of Cluster objects.
*/
public Cluster[] getClusters();
}
7- KMeansFrame
package kmeans;
import java.awt.*;
import javax.swing.*;
import java.awt.event.*;
import java.io.*;
import java.util.*;
import java.sql.*;
import javax.swing.JOptionPane;
/**
* Swing-based program for testing the versions of K-Means on
* randomly generated data.
*/
public class KMeansFrame extends JFrame
implements ActionListener,
KMeansListener {
JPanel contentPane;
BorderLayout borderLayout1 = new BorderLayout();
JScrollPane mMessageAreaSP = new JScrollPane();
JTextArea mMessageArea = new JTextArea();
JButton mRunButton = new JButton();
JPanel mTopPanel = new JPanel();
JLabel mImplementationLabel = new JLabel();
JTextField mRandomSeedTF = new JTextField();
JLabel mRandomSeedLabel = new JLabel();
JTextField mClusterCountTF = new JTextField();
JLabel mClusterCountLabel = new JLabel();
JTextField mCoordCountTF = new JTextField();
JLabel mCountLabel = new JLabel();
JComboBox mImplementationCB = new JComboBox();
GridBagLayout gridBagLayout1 = new GridBagLayout();
JLabel mThreadCountLabel = new JLabel();
JTextField mThreadCountTF = new JTextField();
private boolean mRunning;
private KMeans mKMeans;
private static final String BASIC_KMEANS = "Basic K-Means Clustering";
private static final String BENCHMARKED_KMEANS = "Benchmarked K-Means Clustering";
private static final String CONCURRENT_KMEANS = "Concurrent K-Means Clustering";
public KMeansFrame() {
setDefaultCloseOperation(EXIT_ON_CLOSE);
contentPane = (JPanel) getContentPane();
contentPane.setLayout(borderLayout1);
setSize(new Dimension(620, 560));
setTitle("K-MEANS by Ahmad Traboulsi");
Font font = new Font("Bookman Old Style", Font.BOLD, 14);
mMessageArea.setFont(font);
mMessageArea.setForeground(Color.BLUE);
mMessageArea.setText("");
mRunButton.setText("Run KMeans");
mRunButton.addActionListener(this);
mImplementationLabel.setText("KMeans Implementation:");
mImplementationCB.addActionListener(this);
mRandomSeedTF.setText("1234");
mRandomSeedTF.setColumns(10);
mRandomSeedLabel.setText("Random Seed:");
mClusterCountTF.setText("30");
mClusterCountTF.setColumns(10);
mClusterCountLabel.setText("Number of Clusters (K):");
mCoordCountTF.setText("2500");
mCoordCountTF.setColumns(10);
mCountLabel.setText("Number of Coordinates (N):");
mThreadCountLabel.setEnabled(false);
mThreadCountLabel.setText("Number of Threads:");
mThreadCountTF.setEnabled(false);
// Initialize the thread count textfield with the
// number of available processors.
mThreadCountTF.setText(String.valueOf(Runtime.getRuntime().availableProcessors()));
mThreadCountTF.setColumns(10);
mTopPanel.setLayout(gridBagLayout1);
contentPane.add(mMessageAreaSP, java.awt.BorderLayout.CENTER);
contentPane.add(mTopPanel, java.awt.BorderLayout.NORTH);
mMessageAreaSP.getViewport().add(mMessageArea);
mTopPanel.add(mCountLabel, new GridBagConstraints(0, 0, 1, 1, 0.0, 0.0
, GridBagConstraints.EAST, GridBagConstraints.NONE,
new Insets(10, 10, 5, 0), 0, 0));
mTopPanel.add(mCoordCountTF,
new GridBagConstraints(1, 0, 1, 1, 0.0, 0.0
, GridBagConstraints.WEST,
GridBagConstraints.NONE,
new Insets(10, 0, 5, 10), 0, 0));
mTopPanel.add(mClusterCountLabel,
new GridBagConstraints(0, 1, 1, 1, 0.0, 0.0
, GridBagConstraints.EAST,
GridBagConstraints.NONE,
new Insets(0, 10, 5, 0), 0, 0));
mTopPanel.add(mClusterCountTF,
new GridBagConstraints(1, 1, 1, 1, 0.0, 0.0
, GridBagConstraints.WEST,
GridBagConstraints.NONE,
new Insets(0, 0, 5, 10), 0, 0));
mTopPanel.add(mRandomSeedLabel,
new GridBagConstraints(0, 2, 1, 1, 0.0, 0.0
, GridBagConstraints.EAST,
GridBagConstraints.NONE,
new Insets(0, 10, 5, 0), 0, 0));
mTopPanel.add(mRandomSeedTF,
new GridBagConstraints(1, 2, 1, 1, 0.0, 0.0
, GridBagConstraints.WEST,
GridBagConstraints.NONE,
new Insets(0, 0, 5, 10), 0, 0));
mTopPanel.add(mImplementationLabel,
new GridBagConstraints(2, 0, 1, 1, 0.0, 0.0
, GridBagConstraints.WEST,
GridBagConstraints.NONE,
new Insets(0, 10, 5, 0), 0, 0));
mTopPanel.add(mImplementationCB,
new GridBagConstraints(3, 0, 1, 1, 1.0, 0.0
, GridBagConstraints.CENTER,
GridBagConstraints.HORIZONTAL,
new Insets(0, 2, 5, 10), 0, 0));
mTopPanel.add(mThreadCountLabel,
new GridBagConstraints(2, 1, 1, 1, 0.0, 0.0
, GridBagConstraints.EAST,
GridBagConstraints.NONE,
new Insets(0, 10, 5, 0), 0, 0));
mTopPanel.add(mThreadCountTF,
new GridBagConstraints(3, 1, 1, 1, 0.0, 0.0
, GridBagConstraints.WEST,
GridBagConstraints.NONE,
new Insets(0, 2, 5, 10), 0, 0));
mTopPanel.add(mRunButton, new GridBagConstraints(3, 2, 1, 1, 0.0, 0.0
, GridBagConstraints.CENTER, GridBagConstraints.HORIZONTAL,
new Insets(0, 0, 5, 10), 0, 0));
mImplementationCB.addItem(BASIC_KMEANS);
mImplementationCB.addItem(BENCHMARKED_KMEANS);
mImplementationCB.addItem(CONCURRENT_KMEANS);
}
/**
* Method for validating entries typed into text fields.
*/
private static long getEnteredValue(JTextField tf, long min, long max) {
long value = 0L;
String s = tf.getText().trim();
if (s.length() == 0) {
throw new RuntimeException("blank entry");
}
try {
value = Long.parseLong(s);
if (value < min || value > max) {
throw new RuntimeException("not in range [" + min + " - " + max + "]");
}
} catch (NumberFormatException nfe) {
throw new RuntimeException("invalid number");
}
return value;
}
/**
* Generates the coordinates to be clustered.
*
* @param coordCount the number of coordinates.
* @param dimensions the length of the coordinates.
* @param clusterCount the number of clusters in the distribution.
* @param randomSeed the seed used by the random number generator.
* @return
*/
private static double[][] generateCoordinates( int coordCount, int dimensions, int clusterCount, long randomSeed)
throws InsufficientMemoryException {
// Explicit garbage collection to reduce the likelihood of
// having insufficient memory.
System.gc();
// JOptionPane.showMessageDialog(null,""+"free memory = "+Runtime.getRuntime().freeMemory()+" Required memory"+( 8L * (long) dimensions * (long) (coordCount + clusterCount)));
long memRequired = 8L * (long) dimensions * (long) (coordCount + clusterCount);
if (Runtime.getRuntime().freeMemory() < memRequired) {
throw new InsufficientMemoryException();
}
double[][] coordinates = new double[coordCount][dimensions];
double[][] exemplars = new double[clusterCount][dimensions];
int response = JOptionPane.showConfirmDialog(null, "Do you want to Create New values?", "Confirm",
JOptionPane.YES_NO_OPTION, JOptionPane.QUESTION_MESSAGE);
// JOptionPane.showMessageDialog(null,response);
if (response == JOptionPane.NO_OPTION)
{
//Read coordinates from database
try {
// Step 1: Load the JDBC driver.
Class.forName("sun.jdbc.odbc.JdbcOdbcDriver");
// Step 2: Establish the connection to the database.
String url = "jdbc:odbc:SQLServerDS";
//Get the database numbers
Connection conn2 = DriverManager.getConnection(url,"sa","");
Statement stmt = conn2.createStatement();
String query2 = "select * from Table3";
ResultSet rs = stmt.executeQuery(query2);
int i=0, j=0;
int y=0;
String x="";
//save the db numbers into a double dimentional array
while (rs.next())
{
while(i< coordCount && j<dimensions)
{
y = j+1;
x = "num"+y;
coordinates[j]=Double.parseDouble(rs.getString(x));
j++;
}
j=0;
i++;
}
stmt.close();
conn2.close();
} catch (Exception e) {
System.err.println("Got an exception!");
System.err.println(e.getMessage()); }
}
else
if (response == JOptionPane.YES_OPTION)
{
// Random coordinates
// double[][] coordinates = new double[coordCount][dimensions];
// double[][] exemplars = new double[clusterCount][dimensions];
Random random = new Random(randomSeed);
for (int i=0; i<clusterCount; i++) {
for (int j=0; j<dimensions; j++) {
exemplars[j] = 100.0 * random.nextDouble();
}
}
for (int i=0; i<coordCount; i++) {
int cluster = random.nextInt(clusterCount);
double[] exemplar = exemplars[cluster];
double[] coord = coordinates;
for (int j=0; j<dimensions; j++) {
coord[j] = exemplar[j] + 50*random.nextGaussian();
}
}
//Write the coordinates into the database
// JOptionPane.showMessageDialog(null,coordCount+" "+dimensions);
try {
// Step 1: Load the JDBC driver.
Class.forName("sun.jdbc.odbc.JdbcOdbcDriver");
// Step 2: Establish the connection to the database.
String url = "jdbc:odbc:SQLServerDS";
Connection conn5 = DriverManager.getConnection(url,"sa","");
Statement stmt5 = conn5.createStatement();
//create the table that will hold the double dimensional array of coordinates.
String query5="drop table Table3;create table Table3 (";
for (int u=1; u <= dimensions; u++)
query5 += "num"+u+" float null, ";
query5 += ");";
query5 = query5.replace(", );" , ");");
stmt5.executeUpdate(query5);
stmt5.close();
conn5.close();
Connection conn4 = DriverManager.getConnection(url,"sa","");
Statement stmt4 = conn4.createStatement();
for (int ii=0; ii< coordCount; ii++)
{
String query4="insert into Table3 values (";
for (int jj=0; jj< dimensions; jj++)
query4 += ""+coordinates[ii][jj]+", ";
query4 += ");";
query4 = query4.replace(", );" , ");");
stmt4.executeUpdate(query4);
}
stmt4.close();
conn4.close();
}
catch (Exception e) {
System.err.println("Got an exception! ");
System.err.println(e.getMessage()); }
}
return coordinates;
}
public synchronized void actionPerformed(ActionEvent e) {
if (e.getSource() == mRunButton && !mRunning) {
// Ensure entered parameters make sense.
try {
int coordCount = (int) getEnteredValue (mCoordCountTF,
1L, (long) Integer.MAX_VALUE);
int clusterCount = (int) getEnteredValue(mClusterCountTF,
1L, (long) (coordCount - 1));
long randomSeed = getEnteredValue(mRandomSeedTF,
Long.MIN_VALUE, Long.MAX_VALUE);
// 5 is the dimensions number
double[][] coordinates = generateCoordinates(coordCount, 5, clusterCount, randomSeed);
String implementation = (String) mImplementationCB.getSelectedItem();
if (implementation == BASIC_KMEANS) {
mKMeans = new BasicKMeans(coordinates, clusterCount, 500, randomSeed);
} else if (implementation == BENCHMARKED_KMEANS) {
mKMeans = new BenchmarkedKMeans(coordinates, clusterCount, 500, randomSeed);
} else if (implementation == CONCURRENT_KMEANS) {
try {
int threadCount = (int) getEnteredValue(mThreadCountTF, 1L, 20L);
mKMeans = new ConcurrentKMeans(coordinates, clusterCount, 500, randomSeed, threadCount);
} catch (RuntimeException rte2) {
JOptionPane.showMessageDialog(this,
"The thread count entry is invalid (" + rte2.getMessage() +
").\nPlease enter a thread count in the range [1 - 20].",
"Invalid Entry", JOptionPane.ERROR_MESSAGE);
return;
}
}
if (mKMeans != null) {
// Force gc, so resources taken up by a previous run that haven't
// been freed yet, will not affect this test.
System.gc();
mMessageArea.setText("");
mKMeans.addKMeansListener(this);
mRunButton.setEnabled(false);
new Thread(mKMeans).start();
mRunning = true;
}
} catch (InsufficientMemoryException ime) {
displayInsufficientMemoryDialog();
} catch (RuntimeException rte) {
JOptionPane.showMessageDialog(this,
"One or more entries are invalid (" + rte.getMessage() +
").\nPlease enter positive numbers for the number of coordinates and clusters\n" +
"The number of clusters must be less than the number of coordinates.",
"Invalid Entries", JOptionPane.ERROR_MESSAGE);
}
} else if (e.getSource() == mImplementationCB) {
boolean b = mImplementationCB.getSelectedItem() == CONCURRENT_KMEANS;
mThreadCountLabel.setEnabled(b);
mThreadCountTF.setEnabled(b);
}
}
/**
* Displays an error dialog stating that insufficient memory is
* available.
*/
private void displayInsufficientMemoryDialog() {
JOptionPane.showMessageDialog(this,
"Insufficient memory is available. Try reducing the \n" +
"number of coordinates and/or the number of clusters.",
"Insufficient Memory", JOptionPane.ERROR_MESSAGE);
}
/**
* Cleanup after completion of k-means.
*/
private void cleanupAfterKMeans() {
if (mKMeans != null) {
mKMeans.removeKMeansListener(this);
mKMeans = null;
}
mRunning = false;
mRunButton.setEnabled(true);
}
public void kmeansMessage(String message) {
displayText(message);
}
public void kmeansComplete(Cluster[] clusters, long executionTime) {
displayText("K-Means complete: processing time (ms) = " + executionTime);
displayText("Number of clusters: " + clusters.length);
cleanupAfterKMeans();
}
public void kmeansError(Throwable err) {
cleanupAfterKMeans();
if (err instanceof InsufficientMemoryException) {
displayText("K-Means aborted because of insufficient memory.");
displayInsufficientMemoryDialog();
} else {
Throwable t = (Throwable) err;
StringWriter sw = new StringWriter();
PrintWriter pw = new PrintWriter(sw);
t.printStackTrace(pw);
displayText(sw.toString());
}
}
private void displayText(String text) {
mMessageArea.append(text);
if (!text.endsWith("\n")) {
mMessageArea.append("\n");
}
}
/**
* Application entry point.
*
* @param args String[]
*/
public static void main(String[] args) {
SwingUtilities.invokeLater(new Runnable() {
public void run() {
try {
UIManager.setLookAndFeel(UIManager.
getSystemLookAndFeelClassName());
} catch (Exception exception) {
exception.printStackTrace();
}
KMeansFrame frame = new KMeansFrame();
frame.validate();
// Center the window
Dimension screenSize = Toolkit.getDefaultToolkit().getScreenSize();
Dimension frameSize = frame.getSize();
if (frameSize.height > screenSize.height) {
frameSize.height = screenSize.height;
}
if (frameSize.width > screenSize.width) {
frameSize.width = screenSize.width;
}
frame.setLocation((screenSize.width - frameSize.width) / 2,
(screenSize.height - frameSize.height) / 2);
frame.setVisible(true);
}
});
}
}
8- KMeansListener
package kmeans;
/**
* Defines object which register with implementation of <code>KMeans</code>
* to be notified of significant events during clustering.
*/
public interface KMeansListener {
/**
* A message has been received.
*
* @param message
*/
public void kmeansMessage(String message);
/**
* KMeans is complete.
*
* @param clusters the output of clustering.
* @param executionTime the time in milliseconds taken to cluster.
*/
public void kmeansComplete(Cluster[] clusters, long executionTime);
/**
* An error occurred during KMeans clustering.
*
* @param t
*/
public void kmeansError(Throwable t);
}