Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion build.sbt
Original file line number Diff line number Diff line change
Expand Up @@ -280,7 +280,8 @@ ThisBuild / mimaBinaryIssueFilters ++= {
ProblemFilters.exclude[MissingClassProblem]("fs2.kafka.KafkaConsumer$AssignmentSignals$EagerSignals"),
ProblemFilters.exclude[MissingClassProblem]("fs2.kafka.KafkaConsumer$AssignmentSignals$EagerSignals$"),
ProblemFilters.exclude[MissingClassProblem]("fs2.kafka.KafkaConsumer$AssignmentSignals$GracefulSignals"),
ProblemFilters.exclude[MissingClassProblem]("fs2.kafka.KafkaConsumer$AssignmentSignals$GracefulSignals$")
ProblemFilters.exclude[MissingClassProblem]("fs2.kafka.KafkaConsumer$AssignmentSignals$GracefulSignals$"),
ProblemFilters.exclude[ReversedMissingMethodProblem]("fs2.kafka.KafkaCommitter.source")
)
// scalafmt: {}
}
Expand Down
44 changes: 41 additions & 3 deletions modules/core/src/main/scala/fs2/kafka/KafkaCommitter.scala
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,30 @@ sealed abstract class KafkaCommitter[F[_]] { self =>
/**
* Creates a new [[KafkaCommitter]] in which the effect type has been changed using the specified
* `FunctionK`.
*
* The resulting committer is equal to this one, so offsets committed through either of them can
* still be merged into a single commit by [[CommittableOffsetBatch]].
*/
final def mapK[G[_]](f: F ~> G): KafkaCommitter[G] =
KafkaCommitter(offsets => f(self.commit(offsets)), f(self.metadata))
KafkaCommitter.create(source, offsets => f(self.commit(offsets)), f(self.metadata))

/**
* The committer this one was derived from through [[mapK]], or this committer itself. Committers
* with the same source commit through the same consumer, so this is what equality is based on.
*/
private[kafka] def source: AnyRef

final override def equals(that: Any): Boolean =
that match {
case that: KafkaCommitter[_] => source eq that.source
case _ => false
}

final override def hashCode(): Int =
System.identityHashCode(source)

final override def toString: String =
"KafkaCommitter$" + System.identityHashCode(source)

}

Expand All @@ -51,8 +72,25 @@ object KafkaCommitter {
override val metadata: F[ConsumerGroupMetadata] =
consumerGroupMetadata

override def toString: String =
"KafkaCommitter$" + System.identityHashCode(this)
override val source: AnyRef =
this

}

private def create[F[_]](
committerSource: AnyRef,
commitOffsets: Map[TopicPartition, OffsetAndMetadata] => F[Unit],
consumerGroupMetadata: F[ConsumerGroupMetadata]
): KafkaCommitter[F] =
new KafkaCommitter[F] {
override def commit(offsets: Map[TopicPartition, OffsetAndMetadata]): F[Unit] =
commitOffsets(offsets)

override val metadata: F[ConsumerGroupMetadata] =
consumerGroupMetadata

override val source: AnyRef =
committerSource

}

Expand Down
38 changes: 38 additions & 0 deletions modules/core/src/test/scala/fs2/kafka/CommittableOffsetSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ package fs2.kafka

import cats.~>
import cats.data.OptionT
import cats.effect.unsafe.implicits.global
import cats.effect.IO
import cats.effect.SyncIO

import org.apache.kafka.clients.consumer.OffsetAndMetadata
Expand Down Expand Up @@ -62,5 +64,41 @@ final class CommittableOffsetSpec extends BaseSpec {

assert(committed == Map(partition -> offsetAndMetadata))
}

it("should keep offsets from the same committer batchable after mapK") {
val partition0 = new TopicPartition("topic", 0)
val partition1 = new TopicPartition("topic", 1)
var committed: List[Map[TopicPartition, OffsetAndMetadata]] = Nil

val committer =
KafkaCommitter[IO](
offsets => IO { committed = offsets :: committed },
IO.raiseError(new NotImplementedError)
)

val f = new (IO ~> OptionT[IO, *]) {
override def apply[A](fa: IO[A]): OptionT[IO, A] = OptionT.liftF(fa)
}

val offsets =
List(
CommittableOffset[IO](partition0, new OffsetAndMetadata(1L), committer),
CommittableOffset[IO](partition1, new OffsetAndMetadata(5L), committer),
CommittableOffset[IO](partition0, new OffsetAndMetadata(2L), committer)
).map(_.mapK(f))

assert(offsets.map(_.committer).distinct.size == 1)

val batch = CommittableOffsetBatch.fromFoldable(offsets)

batch.commit.value.unsafeRunSync()

assert(batch.offsets.size == 1)
assert(
committed == List(
Map(partition0 -> new OffsetAndMetadata(2L), partition1 -> new OffsetAndMetadata(5L))
)
)
}
}
}
Loading