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

Fix dataloss issue in restarting-streams-at-rebalancing mode #473

Merged
merged 12 commits into from Sep 23, 2022
@@ -0,0 +1,18 @@
package zio.kafka.consumer.internal

import zio._
import zio.kafka.consumer.internal.Runloop.ByteArrayCommittableRecord
import zio.stream.Take

private[internal] case class PartitionStreamControl(
interrupt: Promise[Throwable, Unit],
drainQueue: Queue[Take[Nothing, ByteArrayCommittableRecord]]
) {

def finishWith(remaining: Chunk[ByteArrayCommittableRecord]): ZIO[Any, Nothing, Unit] =
for {
_ <- drainQueue.offer(Take.chunk(remaining))
_ <- drainQueue.offer(Take.end)
_ <- interrupt.succeed(())
} yield ()
}