From 5d1f8fc4e2e60d6c67d86af49daae41735a49d6b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alejandro=20Herna=CC=81ndez?= Date: Wed, 30 Sep 2026 16:00:45 +0200 Subject: [PATCH] Lift records with a plain mapK on fs2-kafka 4.1.2 fs2-kafka 4.1.2 (typelevel/fs2-kafka#1531) makes a committer produced by `KafkaCommitter#mapK` equal to the committer it was derived from, so `CommittableOffsetBatch` merges lifted offsets again. `injectK` no longer needs to cache one lifted committer per consumer, and the `fs2.kafka.LiftedCommittableOffset` shim that rebuilt offsets around the cached committer goes with it. 4.1.2 also fixes the deadlock on `KafkaConsumer#unsubscribe` once partitions are assigned (typelevel/fs2-kafka#1530). --- .../fs2/kafka/LiftedCommittableOffset.scala | 34 ------------------- .../trace4cats/kafka/TracedConsumer.scala | 23 ++----------- project/dependencies.conf | 2 +- 3 files changed, 3 insertions(+), 56 deletions(-) delete mode 100644 modules/trace4cats-kafka-client/src/main/scala/fs2/kafka/LiftedCommittableOffset.scala diff --git a/modules/trace4cats-kafka-client/src/main/scala/fs2/kafka/LiftedCommittableOffset.scala b/modules/trace4cats-kafka-client/src/main/scala/fs2/kafka/LiftedCommittableOffset.scala deleted file mode 100644 index 5b84402..0000000 --- a/modules/trace4cats-kafka-client/src/main/scala/fs2/kafka/LiftedCommittableOffset.scala +++ /dev/null @@ -1,34 +0,0 @@ -/* - * Copyright (c) 2021-2026 Trace4Cats - * - * Permission is hereby granted, free of charge, to any person obtaining a copy of - * this software and associated documentation files (the "Software"), to deal in - * the Software without restriction, including without limitation the rights to - * use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of - * the Software, and to permit persons to whom the Software is furnished to do so, - * subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS - * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR - * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER - * IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN - * CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. - */ - -package fs2.kafka - -/** Rebuilds a [[CommittableOffset]] around a committer already lifted to another effect. Lives in `fs2.kafka` because - * `CommittableOffset.apply` is package-private; the alternative, `CommittableOffset#mapK`, lifts the committer anew on - * every call and [[CommittableOffsetBatch]] groups offsets by committer identity, so records lifted that way never - * batch their commits together. - */ -object LiftedCommittableOffset { - - def apply[F[_], G[_]](offset: CommittableOffset[F], committer: KafkaCommitter[G]): CommittableOffset[G] = - CommittableOffset(offset.topicPartition, offset.offsetAndMetadata, committer) - -} diff --git a/modules/trace4cats-kafka-client/src/main/scala/trace4cats/kafka/TracedConsumer.scala b/modules/trace4cats-kafka-client/src/main/scala/trace4cats/kafka/TracedConsumer.scala index 4373ee8..5c46757 100644 --- a/modules/trace4cats-kafka-client/src/main/scala/trace4cats/kafka/TracedConsumer.scala +++ b/modules/trace4cats-kafka-client/src/main/scala/trace4cats/kafka/TracedConsumer.scala @@ -22,7 +22,6 @@ package trace4cats.kafka import cats.Functor -import cats.data.WriterT import cats.effect.kernel.MonadCancelThrow import cats.syntax.applicativeError._ import cats.syntax.functor._ @@ -30,7 +29,6 @@ import cats.syntax.functor._ import fs2.Stream import fs2.kafka.CommittableConsumerRecord import fs2.kafka.KafkaCommitter -import fs2.kafka.LiftedCommittableOffset import fs2.kafka.Timestamp import trace4cats.ResourceKleisli import trace4cats.Span @@ -88,24 +86,7 @@ object TracedConsumer extends Fs2StreamSyntax { stream: Stream[F, CommittableConsumerRecord[F, K, V]] )( k: ResourceKleisli[F, SpanParams, Span[F]] - )(implicit P: Provide[F, G, Span[F]]): TracedStream[G, CommittableConsumerRecord[G, K, V]] = { - val liftK = P.liftK - - WriterT( - inject[F, G, K, V](stream)(k) - .liftTrace[G] - .run - .mapAccumulate(Map.empty[KafkaCommitter[F], KafkaCommitter[G]]) { case (committers, (span, record)) => - val committer = committers.getOrElse(record.offset.committer, record.offset.committer.mapK(liftK)) - val offset = LiftedCommittableOffset(record.offset, committer) - - ( - committers.updated(record.offset.committer, committer), - (span, CommittableConsumerRecord(record.record, offset)) - ) - } - .map(_._2) - ) - } + )(implicit P: Provide[F, G, Span[F]]): TracedStream[G, CommittableConsumerRecord[G, K, V]] = + inject[F, G, K, V](stream)(k).liftTrace[G].map(_.mapK(P.liftK)) } diff --git a/project/dependencies.conf b/project/dependencies.conf index 87eaead..68f6c1e 100644 --- a/project/dependencies.conf +++ b/project/dependencies.conf @@ -23,6 +23,6 @@ common-settings { trace4cats-kafka-client = [ "io.janstenpickle::trace4cats-fs2:0.14.7" - "org.typelevel::fs2-kafka:4.1.1" + "org.typelevel::fs2-kafka:4.1.2" "io.janstenpickle::trace4cats-testkit:0.14.7:test" ]