fixing memory leak in kafka MessageHandler

This commit is contained in:
seanm 2013-03-14 23:25:35 -06:00
parent cfa8e769a8
commit d069283211

View file

@ -114,11 +114,8 @@ class KafkaReceiver(kafkaParams: Map[String, String],
private class MessageHandler(stream: KafkaStream[String]) extends Runnable {
def run() {
logInfo("Starting MessageHandler.")
stream.takeWhile { msgAndMetadata =>
for (msgAndMetadata <- stream) {
blockGenerator += msgAndMetadata.message
// Keep on handling messages
true
}
}
}