First, is you need to tell Kafka the name of your Partitioner class:
props.put("partitioner.class", "com.curtinhome.kafka.playproducer.OrganizationPartitioner");
Then tell Kafka how to serialize the key so the Partitioner class can be called with the value:
props.put("serializer.class", "kafka.serializer.StringEncoder");
props.put("key.serializer.class", "kafka.serializer.StringEncoder");
Here we're telling Kafka that both the message being written and the key are Strings. If you wanted to pass an Integer key, you'd need to provide a class that coverts the Integer to a byte array.
Now when you write the message, you pass the Key and the Message:
long events = Long.parseLong(args[0]);
int blocks = Integer.parseInt(args[1]);
Random rnd = new Random();
Properties props = new Properties();
props.put("broker.list", "broker1.atlnp1:9092,broker2.atlnp1:9092,broker3.atlnp1:9092");
props.put("serializer.class", "kafka.serializer.StringEncoder");
props.put("key.serializer.class", "kafka.serializer.StringEncoder");
props.put("partitioner.class", "com.curtinhome.kafka.playproducer.OrganizationPartitioner");
ProducerConfig config = new ProducerConfig(props);
Producer
for (int nBlocks = 0; nBlocks < blocks; nBlocks++) {
for (long nEvents = 0; nEvents < events; nEvents++) {
long runtime = new Date().getTime();
String msg = runtime + "," + (50 + nBlocks) + "," + nEvents+ "," + rnd.nextInt(1000);
KeyedMessage
new KeyedMessage
producer.send(data);
}
}
producer.close();
The Partitioner is pretty straight forward:
public class OrganizationPartitioner implements Partitioner
public OrganizationPartitioner(VerifiableProperties props) {
}
public int partition(String key, int a_numPartitions) {
long organizationId = Long.parseLong(key);
return (int) (organizationId % a_numPartitions);
}
}
The number of partitions is provided by Kafka based on the topic configuration.