-
Notifications
You must be signed in to change notification settings - Fork 62
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Question : I try to use it but nothing ... #26
Comments
Hmm, perhaps there is an off by one error. What happens if you use
also, is there any process sending messages to kafka? |
also, what version are you using? |
Thanks for super quick help !! using 0.9 from maven repo here is my log with debug |
OOh, sorry I think 0.0.9 is broken. .Can you try 0.0.10-SNAPSHOT from On Tue, Jan 13, 2015 at 9:23 AM, richiesgr notifications@github.com wrote:
|
0.0.10 should be live in just a few minutes On Tue, Jan 13, 2015 at 9:24 AM, Scott Clasen scott@heroku.com wrote:
|
Thanks I give it a try |
Thanks a lot for your help working now on 0.0.10 |
Hi
I've try to use it but I can't get it working I mean that I get message from kafka only at the this time then I enter in this mode:
[test-akka.actor.default-dispatcher-3] INFO com.sclasen.akka.kafka.ConnectorFSM - at=created-streams
[test-akka.actor.default-dispatcher-3] INFO akka.actor.ActorSystemImpl - at=consumer-started
[test-akka.actor.default-dispatcher-5] INFO com.sclasen.akka.kafka.ConnectorFSM - at=transition from=Receiving to=Committing uncommitted=0
[test-akka.actor.default-dispatcher-5] WARN com.sclasen.akka.kafka.ConnectorFSM - state=Committing msg=StateTimeout drained=1 streams=1
[test-akka.actor.default-dispatcher-5] WARN com.sclasen.akka.kafka.ConnectorFSM - state=Committing msg=StateTimeout drained=1 streams=1
[test-akka.actor.default-dispatcher-2] WARN com.sclasen.akka.kafka.ConnectorFSM - state=Committing msg=StateTimeout drained=1 streams=1
And don't get nothing anymore.
I've try to play with CommitConf without success.
Could you help me
My code is
Main () {
val system = ActorSystem("test")
val senderRequest = system.actorOf(Props[SenderRequest])
val consumerProps = AkkaConsumerProps.forSystem(
system = system,
zkConnect = "localhost:2181",
topic = "HTTP",
group = "1",
streams = 1, //one per partition
keyDecoder = new DefaultDecoder(),
msgDecoder = new DefaultDecoder(),
receiver = senderRequest,
commitConfig = new CommitConfig(Some(10 seconds), Some(1), Timeout(10 second))
)
val consumer = new AkkaConsumer(consumerProps)
consumer.start()
}
class SenderRequest extends Actor with ActorLogging {
override def receive: Actor.Receive = {
case data : Any =>
log.info("RECEIVE HERE !!!!!!!!!!!!!!!")
sender() ! StreamFSM.Processed
}
}
Thanks
The text was updated successfully, but these errors were encountered: