final int CONCURRENCY = 5;
return infiniteSubmitter.start() //
.flatMap(notUsed -> getFromMessageRouter(getDmaapUrl()), 1) //
final int CONCURRENCY = 5;
return infiniteSubmitter.start() //
.flatMap(notUsed -> getFromMessageRouter(getDmaapUrl()), 1) //
.doOnNext(message -> logger.debug("Message Reveived from DMAAP : {}", message)) //
.flatMap(this::handleReceivedMessage, CONCURRENCY);
}
.doOnNext(message -> logger.debug("Message Reveived from DMAAP : {}", message)) //
.flatMap(this::handleReceivedMessage, CONCURRENCY);
}