buggood first issue
Metriche repository
- Star
- (124 stelle)
- Metriche merge PR
- (Metriche PR in attesa)
Descrizione
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