From 5bff613ee3b4eab0a7696f73e32ede2975fadb1f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alejandro=20Herna=CC=81ndez?= Date: Tue, 29 Sep 2026 19:04:38 +0200 Subject: [PATCH 1/2] Ensure lifted committers carry a reference to root committer so they're not duplicated New `mapK` function introduced in #1522 creates a new committer on every call and its `Eq` instance is based on `equals` so every mapK generates a new committer equal to nothing else. The problem is CommittableOffsetBatch groups by committer, so a lifted stream like the one we create in trace4cats-kafka fell back to one comit per record, in parallel. This solves it by adding a "pointer" to the source committer, so that can be used for `hashCode` and `equals`. --- .../main/scala/fs2/kafka/KafkaCommitter.scala | 44 +++++++++++++++++-- .../fs2/kafka/CommittableOffsetSpec.scala | 38 ++++++++++++++++ 2 files changed, 79 insertions(+), 3 deletions(-) diff --git a/modules/core/src/main/scala/fs2/kafka/KafkaCommitter.scala b/modules/core/src/main/scala/fs2/kafka/KafkaCommitter.scala index 315eb909a..29e07341d 100644 --- a/modules/core/src/main/scala/fs2/kafka/KafkaCommitter.scala +++ b/modules/core/src/main/scala/fs2/kafka/KafkaCommitter.scala @@ -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) } @@ -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 } diff --git a/modules/core/src/test/scala/fs2/kafka/CommittableOffsetSpec.scala b/modules/core/src/test/scala/fs2/kafka/CommittableOffsetSpec.scala index c651e53d6..288d2fb72 100644 --- a/modules/core/src/test/scala/fs2/kafka/CommittableOffsetSpec.scala +++ b/modules/core/src/test/scala/fs2/kafka/CommittableOffsetSpec.scala @@ -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 @@ -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)) + ) + ) + } } } From 08fac6447e351eb51b2e6f6303a5266dee6c0575 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alejandro=20Herna=CC=81ndez?= Date: Tue, 29 Sep 2026 19:05:51 +0200 Subject: [PATCH 2/2] Add filter for new KafkaCommitter#source The class is sealed with a package-private constructor so no one outside this repo can implement it. --- build.sbt | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/build.sbt b/build.sbt index 88b34520b..be9f5ca31 100644 --- a/build.sbt +++ b/build.sbt @@ -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: {} }