How do I distribute keys into buckets but ensure that each bucket has at least one key?

Viewed 203

I have a set of "keys" and "buckets", where num_keys > num_buckets, but they are roughly the same order of magnitude.

To give a concrete example, a "key" would be analogous to a task, and a "bucket" to a thread. I want to distribute the tasks as evenly as possible to the threads, while ensuring that no thread is idle.

For example, I would have 400 tasks, and 200 threads:

task1, task2, ..., task400
thread1, thread1, ... thread200

And I want to assign tasks somewhat uniformly to the threads:

thread1 -> task50, task394
thread2 -> task37, task250, task324
...
thread200 -> task20

These are numbers I made up, but to be clear, it's more important to me that there are no idle threads; I don't care that much about having perfect distribution.

Some other properties that are important to me:

  • There is no bias toward a particular bucket.
  • The same key should always be assigned to the same bucket, regardless of the number of keys. If the number of buckets change, it's OK for the assignments to change (so no need for some kind of consistent hashing scheme).
  • The distribution should happen independently without external coordination.

Currently, I'm trying to do this with MD5:

BigInteger numThreads = BigInteger.valueOf(200); // buckets
BigInteger currentThreadId = BigInteger.valueOf(1); 
List<Task> tasks = getGlobalTasks(); // keys
List<Task> myTasks = new ArrayList<>();

for (Task task : tasks) {
    BigInteger hash = new BigInteger(MessageDigest.getInstance("MD5").digest(task.getId()));
    BigInteger bucketId = hash.mod(numThreads);
    if (currentThreadId.equals(bucketId)) {
      myTasks.add(task);
    }
}

I suspect that this method would work fine if there was a greater disparity between the number of keys and buckets. However, with this current scheme, about 35 of my threads are sitting idle because it hasn't been assigned any keys. I'm guessing that this is due to variance from the MD5 hash distribution.

So the question is, is there any way to restructure this code to ensure that each bucket/thread will receive at least one task?

Note: What I'm actually trying to do is assign AWS Kinesis Shards to Apache Storm Spouts in a way that is resilient to resharding (which means changing the number of shards). The details of this aren't important for this question though.

Here is an example of the code sample from above that compiles/runs:

  public static void main(String [] args) throws NoSuchAlgorithmException {
    BigInteger numBuckets= BigInteger.valueOf(20); // buckets
    int numTasks = 40;

    List<BigInteger> tasks = new ArrayList<>(); // keys
    for (int i = 0; i < numTasks; i++) {
      tasks.add(BigInteger.valueOf((int)(Math.random() * 5000)));
    }

    Map<BigInteger, List<BigInteger>> bucketToTasks = new HashMap<>();
    for (BigInteger task : tasks) {
        String taskStr = task.toString();
        BigInteger hash = new BigInteger(MessageDigest.getInstance("MD5").digest(taskStr.getBytes()));
        BigInteger bucketId = hash.mod(numBuckets);
        if (!bucketToTasks.containsKey(bucketId)) {
          bucketToTasks.put(bucketId, new ArrayList<>());
        }
        bucketToTasks.get(bucketId).add(task);
    }

    for (int i = 0; i < numBuckets.intValue(); i++){
      System.out.println("Bucket: " + i + ", Tasks: " + bucketToTasks.get(BigInteger.valueOf(i)));
    }
  }

And the output:

Bucket: 0, Tasks: [2459, 2922]
Bucket: 1, Tasks: [21]
Bucket: 2, Tasks: [453]
Bucket: 3, Tasks: [494]
Bucket: 4, Tasks: null
Bucket: 5, Tasks: [1557, 2355]
Bucket: 6, Tasks: [4145, 1693]
Bucket: 7, Tasks: null
Bucket: 8, Tasks: null
Bucket: 9, Tasks: [1016, 4556]
Bucket: 10, Tasks: [3739]
Bucket: 11, Tasks: [4644, 3295]
Bucket: 12, Tasks: [2773, 3983, 2724]
Bucket: 13, Tasks: [1663]
Bucket: 14, Tasks: [2593, 2593, 1457, 3023, 1817, 4353, 4196]
Bucket: 15, Tasks: [2258, 802, 1530]
Bucket: 16, Tasks: [1882, 4765, 1969]
Bucket: 17, Tasks: [4882, 3030, 3429]
Bucket: 18, Tasks: [1666, 4377, 1714, 4109]
Bucket: 19, Tasks: [819, 7]

As you can see, some buckets are completely empty.

0 Answers
Related