Skip to content
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

Offsets issue #6

Open
ghost opened this issue Apr 26, 2017 · 0 comments
Open

Offsets issue #6

ghost opened this issue Apr 26, 2017 · 0 comments

Comments

@ghost
Copy link

ghost commented Apr 26, 2017

Hi,

I want that each time the inner function in foreachL crashes with exception, the offset to stay the same. It appears it doesn't work like this, changing the consumer group to something else will throw an error first time I run this code(with just one message in topic topic), as expected, but second time will work because foreachL is not called anymore since the record is already consumed.

What am I doing wrong?

Thank You,

object Main extends App {

  private val consumerCfg = KafkaConsumerConfig.default.copy(
    bootstrapServers = List("localhost:9092"),
    groupId = "foo-7",
    enableAutoCommit = false,
    autoOffsetReset = AutoOffsetReset.Earliest,
    observableCommitOrder = ObservableCommitOrder.AfterAck,
    observableCommitType = ObservableCommitType.Sync
  )

  val f = KafkaConsumerObservable[String, String](consumerCfg, List("topic"))
    .timeoutOnSlowUpstream(5.seconds)
    .foreachL { _ ⇒ throw new Exception("crash") }
    .runAsync(Scheduler.io())

  Await.result(f, Duration.Inf)
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment
Labels
None yet
Projects
None yet
Development

No branches or pull requests

0 participants