buggood first issue
Repository metrics
- Stars
- (124 stars)
- PR merge metrics
- (PR metrics pending)
Description
In KafkaCOnsumerObservableManualCommit.scala:56 there's incorrect implementation of async commit
override def commitBatchAsync(batch: Map[TopicPartition, Long], callback: OffsetCommitCallback): Task[Unit] =
Task {
blocking(consumer.synchronized(consumer.commitAsync(batch.map {
case (k, v) => k -> new OffsetAndMetadata(v)
}.asJava, callback)))
}
Apache kafka commitAsync(offsets, callback) returns immediately and invokes callback when commit completes
this task completes when consumer.commitAsync returns, not when callback invoked and also ignores possible commit errors also passed through callback, it should probably be rewritten to Task.async