From 0198a28c9c314b979e0bbc1fac74ed77454fa26b Mon Sep 17 00:00:00 2001 From: Nikita Viliunov Date: Mon, 10 Aug 2026 18:34:29 +0200 Subject: [PATCH] Optimize Converters.topicPartitionsSetF --- .../evolutiongaming/skafka/Converters.scala | 22 ++++++++++++------- .../evolutiongaming/skafka/Partition.scala | 8 ++++--- 2 files changed, 19 insertions(+), 11 deletions(-) diff --git a/skafka/src/main/scala/com/evolutiongaming/skafka/Converters.scala b/skafka/src/main/scala/com/evolutiongaming/skafka/Converters.scala index e98a590f..c03544a4 100644 --- a/skafka/src/main/scala/com/evolutiongaming/skafka/Converters.scala +++ b/skafka/src/main/scala/com/evolutiongaming/skafka/Converters.scala @@ -2,10 +2,9 @@ package com.evolutiongaming.skafka import java.lang.Long as LongJ import java.time.Duration as DurationJ -import java.util.{Optional, Collection as CollectionJ, Map as MapJ, Set as SetJ, List as ListJ} - +import java.util.{Optional, Collection as CollectionJ, List as ListJ, Map as MapJ, Set as SetJ} import cats.Monad -import cats.data.{NonEmptyList as Nel, NonEmptySet as Nes, NonEmptyMap as Nem} +import cats.data.{NonEmptyList as Nel, NonEmptyMap as Nem, NonEmptySet as Nes} import cats.syntax.all.* import com.evolutiongaming.catshelper.CatsHelper.* import com.evolutiongaming.catshelper.{ApplicativeThrowable, FromTry, MonadThrowable, ToTry} @@ -17,6 +16,7 @@ import org.apache.kafka.common.{PartitionInfo as PartitionInfoJ, TopicPartition import scala.compat.java8.DurationConverters import scala.concurrent.duration.FiniteDuration import scala.jdk.CollectionConverters.* +import scala.util.{Failure, Success, Try} object Converters { @@ -233,9 +233,15 @@ object Converters { mapJ.asScalaMap(_.pure[F], partitionsInfoListF[F]) } - def topicPartitionsSetF[F[_]: ApplicativeThrowable](setJ: SetJ[TopicPartitionJ]): F[Set[TopicPartition]] = { - for { - r <- setJ.asScala.toList.traverse { _.asScala[F] } - } yield r.toSet - } + def topicPartitionsSetF[F[_]: ApplicativeThrowable](setJ: SetJ[TopicPartitionJ]): F[Set[TopicPartition]] = + ApplicativeThrowable[F].catchNonFatal { + val builder = Set.newBuilder[TopicPartition] + setJ.forEach { tpj => + tpj.asScala[Try] match { + case Failure(exception) => throw exception + case Success(value) => builder.addOne(value) + } + } + builder.result() + } } diff --git a/skafka/src/main/scala/com/evolutiongaming/skafka/Partition.scala b/skafka/src/main/scala/com/evolutiongaming/skafka/Partition.scala index fad39806..b81a39a4 100644 --- a/skafka/src/main/scala/com/evolutiongaming/skafka/Partition.scala +++ b/skafka/src/main/scala/com/evolutiongaming/skafka/Partition.scala @@ -14,9 +14,11 @@ sealed abstract case class Partition(value: Int) { object Partition { - val min: Partition = new Partition(0) {} + private class Impl(value: Int) extends Partition(value) - val max: Partition = new Partition(Int.MaxValue) {} + val min: Partition = new Impl(0) + + val max: Partition = new Impl(Int.MaxValue) implicit val showPartition: Show[Partition] = Show.fromToString @@ -36,7 +38,7 @@ object Partition { } else if (value == max.value) { max.pure[F] } else { - new Partition(value) {}.pure[F] + (new Impl(value): Partition).pure[F] } }