There are reasons why implementing a system in Java might be a questionable decision, but unless that system involves extremely intensive number-crunching or has hard real-time requirements, performance probably isn't one of them.
http://research.microsoft.com/en-us/um/people/srikanth/netdb...
By "hard real-time" I was referring to latency, rather than throughput. Achieving very low latencies is difficult in Java because you don't know exactly when the GC will kick in, but nevertheless it's possible to get very high throughput.
With Java, you don't have the unmanaged or struct support, so doesn't that really add up? If you go "native", isn't there significant overhead since you can't have pointers in Java (right - the bytecode doesn't support it?)?
People pull it off, but it seems that GC overhead would be a killer.
There's still GC from objects allocated by Kafka in the JVM, but the actual message data doesn't even go through the JVM.