This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors-kafka.git
The following commit(s) were added to refs/heads/main by this push:
new 4052ee33 chore: Rewrite to scala3 syntax (#547)
4052ee33 is described below
commit 4052ee33565d0178fbd5bbbbda2aeefb31bac4e6
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Mon Jun 15 17:13:23 2026 +0800
chore: Rewrite to scala3 syntax (#547)
Motivation:
Let Scala 3 Community build compile pekko-connectors-kafka.
Modification:
Update .scalafmt.conf dialect to scala213source3 and enable
scala3 syntax rewrite rules, then reformat all sources.
Result:
Wildcard types use ? syntax (e.g. Class[?] instead of Class[_]).
Tests:
Not run - formatting only
References:
Refs apache/pekko#3048
---
.scalafmt.conf | 12 +++++++++++-
.../cluster/sharding/KafkaClusterSharding.scala | 18 +++++++++---------
.../apache/pekko/kafka/RestrictedConsumer.scala | 2 +-
.../kafka/internal/CommitCollectorStage.scala | 2 +-
.../internal/CommittingProducerSinkStage.scala | 2 +-
.../pekko/kafka/internal/ConfigSettings.scala | 2 +-
.../kafka/internal/ConsumerProgressTracking.scala | 4 ++--
.../kafka/internal/ControlImplementations.scala | 2 +-
.../kafka/internal/DefaultProducerStage.scala | 4 ++--
.../kafka/internal/ExternalSingleSourceLogic.scala | 2 +-
.../pekko/kafka/internal/KafkaConsumerActor.scala | 8 ++++----
.../pekko/kafka/internal/MessageBuilder.scala | 2 +-
.../pekko/kafka/internal/SingleSourceLogic.scala | 2 +-
.../internal/TransactionalProducerStage.scala | 2 +-
.../kafka/internal/TransactionalSources.scala | 2 +-
.../org/apache/pekko/kafka/javadsl/Producer.scala | 2 +-
.../org/apache/pekko/kafka/scaladsl/Producer.scala | 2 +-
.../org/apache/pekko/kafka/IntegrationTests.scala | 4 ++--
project/CopyrightHeader.scala | 6 +++---
project/CopyrightHeaderForBuild.scala | 2 +-
project/CopyrightHeaderForProtobuf.scala | 2 +-
project/ProjectSettings.scala | 2 +-
project/VersionGenerator.scala | 2 +-
.../KafkaTestkitTestcontainersSettings.scala | 22 +++++++++++-----------
.../kafka/testkit/ProducerResultFactory.scala | 2 +-
.../testkit/internal/TestcontainersKafka.scala | 2 +-
.../testkit/scaladsl/TestcontainersKafkaLike.scala | 2 +-
.../docs/scaladsl/ClusterShardingExample.scala | 2 +-
.../test/scala/docs/scaladsl/ProducerExample.scala | 2 +-
.../src/test/scala/docs/scaladsl/proto/Order.scala | 6 +++---
.../scala/docs/scaladsl/proto/OrderProto.scala | 4 ++--
.../pekko/kafka/internal/ConsumerDummy.scala | 2 +-
.../pekko/kafka/scaladsl/ReconnectSpec.scala | 2 +-
33 files changed, 72 insertions(+), 62 deletions(-)
diff --git a/.scalafmt.conf b/.scalafmt.conf
index e559eb3c..3cd1fbf0 100644
--- a/.scalafmt.conf
+++ b/.scalafmt.conf
@@ -1,5 +1,5 @@
version = 3.11.1
-runner.dialect = scala213
+runner.dialect = scala213source3
project.git = true
style = defaultWithAlign
docstrings.style = Asterisk
@@ -76,3 +76,13 @@ project.excludeFilters = [
"scripts/authors.scala"
]
project.layout = StandardConvention
+
+rewrite.scala3.convertToNewSyntax = true
+runner {
+ dialectOverride {
+ allowSignificantIndentation = false
+ allowAsForImportRename = false
+ allowStarWildcardImport = false
+ allowPostfixStarVarargSplices = false
+ }
+}
diff --git
a/cluster-sharding/src/main/scala/org/apache/pekko/kafka/cluster/sharding/KafkaClusterSharding.scala
b/cluster-sharding/src/main/scala/org/apache/pekko/kafka/cluster/sharding/KafkaClusterSharding.scala
index 7dca330b..cf4069f9 100644
---
a/cluster-sharding/src/main/scala/org/apache/pekko/kafka/cluster/sharding/KafkaClusterSharding.scala
+++
b/cluster-sharding/src/main/scala/org/apache/pekko/kafka/cluster/sharding/KafkaClusterSharding.scala
@@ -65,7 +65,7 @@ final class KafkaClusterSharding(system: ExtendedActorSystem)
extends Extension
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/1074")
def messageExtractor[M](topic: String,
timeout: FiniteDuration,
- settings: ConsumerSettings[_, _]):
Future[KafkaShardingMessageExtractor[M]] =
+ settings: ConsumerSettings[?, ?]):
Future[KafkaShardingMessageExtractor[M]] =
getPartitionCount(topic, timeout, settings).map(new
KafkaShardingMessageExtractor[M](_))(system.dispatcher)
/**
@@ -87,7 +87,7 @@ final class KafkaClusterSharding(system: ExtendedActorSystem)
extends Extension
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/1074")
def messageExtractor[M](topic: String,
timeout: java.time.Duration,
- settings: ConsumerSettings[_, _]):
CompletionStage[KafkaShardingMessageExtractor[M]] =
+ settings: ConsumerSettings[?, ?]):
CompletionStage[KafkaShardingMessageExtractor[M]] =
getPartitionCount(topic, timeout.toScala, settings)
.map(new KafkaShardingMessageExtractor[M](_))(system.dispatcher)
.asJava
@@ -126,7 +126,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
def messageExtractorNoEnvelope[M](topic: String,
timeout: FiniteDuration,
entityIdExtractor: M => String,
- settings: ConsumerSettings[_, _]):
Future[KafkaShardingNoEnvelopeExtractor[M]] =
+ settings: ConsumerSettings[?, ?]):
Future[KafkaShardingNoEnvelopeExtractor[M]] =
getPartitionCount(topic, timeout, settings)
.map(partitions => new KafkaShardingNoEnvelopeExtractor[M](partitions,
entityIdExtractor))(system.dispatcher)
@@ -152,7 +152,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
topic: String,
timeout: java.time.Duration,
entityIdExtractor: java.util.function.Function[M, String],
- settings: ConsumerSettings[_, _]):
CompletionStage[KafkaShardingNoEnvelopeExtractor[M]] =
+ settings: ConsumerSettings[?, ?]):
CompletionStage[KafkaShardingNoEnvelopeExtractor[M]] =
getPartitionCount(topic, timeout.toScala, settings)
.map(partitions => new KafkaShardingNoEnvelopeExtractor[M](partitions, e
=> entityIdExtractor.apply(e)))(
system.dispatcher)
@@ -194,7 +194,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
private val metadataConsumerActorNum = new AtomicInteger
private def getPartitionCount[M](topic: String,
timeout: FiniteDuration,
- settings: ConsumerSettings[_, _]): Future[Int] = {
+ settings: ConsumerSettings[?, ?]): Future[Int] = {
implicit val ec: ExecutionContextExecutor = system.dispatcher
val num = metadataConsumerActorNum.getAndIncrement()
val consumerActor = system
@@ -209,7 +209,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
}
private val rebalanceListeners =
- new ConcurrentHashMap[EntityTypeKey[_],
pekko.actor.typed.ActorRef[ConsumerRebalanceEvent]]()
+ new ConcurrentHashMap[EntityTypeKey[?],
pekko.actor.typed.ActorRef[ConsumerRebalanceEvent]]()
/**
* API MAY CHANGE
@@ -229,7 +229,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
* }}}
*/
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/1074")
- def rebalanceListener(typeKey: EntityTypeKey[_]):
pekko.actor.typed.ActorRef[ConsumerRebalanceEvent] = {
+ def rebalanceListener(typeKey: EntityTypeKey[?]):
pekko.actor.typed.ActorRef[ConsumerRebalanceEvent] = {
rebalanceListeners.computeIfAbsent(typeKey,
_ => {
system.toTyped
@@ -258,7 +258,7 @@ final class KafkaClusterSharding(system:
ExtendedActorSystem) extends Extension
*/
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/1074")
def rebalanceListener(
- typeKey: pekko.cluster.sharding.typed.javadsl.EntityTypeKey[_])
+ typeKey: pekko.cluster.sharding.typed.javadsl.EntityTypeKey[?])
: pekko.actor.typed.ActorRef[ConsumerRebalanceEvent] = {
rebalanceListener(typeKey.asScala)
}
@@ -299,7 +299,7 @@ object KafkaClusterSharding extends
ExtensionId[KafkaClusterSharding] {
// Used from future callbacks so can't use the log in the context
private val log = LoggerFactory.getLogger(RebalanceListener.getClass)
- def apply(typeKey: EntityTypeKey[_]): Behavior[ConsumerRebalanceEvent] =
+ def apply(typeKey: EntityTypeKey[?]): Behavior[ConsumerRebalanceEvent] =
Behaviors.setup { ctx =>
import ctx.executionContext
val shardAllocationClient =
ExternalShardAllocation(ctx.system).clientFor(typeKey.name)
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/RestrictedConsumer.scala
b/core/src/main/scala/org/apache/pekko/kafka/RestrictedConsumer.scala
index 1b5c2f14..5ac180ba 100644
--- a/core/src/main/scala/org/apache/pekko/kafka/RestrictedConsumer.scala
+++ b/core/src/main/scala/org/apache/pekko/kafka/RestrictedConsumer.scala
@@ -24,7 +24,7 @@ import org.apache.kafka.common.TopicPartition
* the [[pekko.kafka.scaladsl.PartitionAssignmentHandler]] callbacks.
*/
@ApiMayChange
-final class RestrictedConsumer(consumer: Consumer[_, _], duration:
java.time.Duration) {
+final class RestrictedConsumer(consumer: Consumer[?, ?], duration:
java.time.Duration) {
/**
* See [[org.apache.kafka.clients.consumer.KafkaConsumer#assignment]]
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/CommitCollectorStage.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/CommitCollectorStage.scala
index 5d471ea4..dd7946f9 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/CommitCollectorStage.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/CommitCollectorStage.scala
@@ -57,7 +57,7 @@ private final class CommitCollectorStageLogic(
private val contextPropagation = pekko.stream.impl.ContextPropagation()
private var contextSuspended = false
- override protected def logSource: Class[_] =
classOf[CommitCollectorStageLogic]
+ override protected def logSource: Class[?] =
classOf[CommitCollectorStageLogic]
private var pushOnNextPull = false
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/CommittingProducerSinkStage.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/CommittingProducerSinkStage.scala
index 78db2033..b2bba846 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/CommittingProducerSinkStage.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/CommittingProducerSinkStage.scala
@@ -71,7 +71,7 @@ private final class CommittingProducerSinkStageLogic[K, V, IN
<: Envelope[K, V,
override protected def getExecutionContext(): ExecutionContext =
materializer.executionContext
- override protected def logSource: Class[_] =
classOf[CommittingProducerSinkStage[_, _, _]]
+ override protected def logSource: Class[?] =
classOf[CommittingProducerSinkStage[?, ?, ?]]
override protected val producerSettings: ProducerSettings[K, V] =
stage.producerSettings
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/ConfigSettings.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/ConfigSettings.scala
index 70ecca8a..bc5c72c5 100644
--- a/core/src/main/scala/org/apache/pekko/kafka/internal/ConfigSettings.scala
+++ b/core/src/main/scala/org/apache/pekko/kafka/internal/ConfigSettings.scala
@@ -38,7 +38,7 @@ import scala.jdk.DurationConverters._
if (unprocessedKeys.isEmpty) processedKeys
else {
c.toConfig.getAnyRef(unprocessedKeys.head) match {
- case o: util.Map[_, _] =>
+ case o: util.Map[?, ?] =>
collectKeys(c,
processedKeys,
unprocessedKeys.tail :::
o.keySet().asScala.toList.map(unprocessedKeys.head + "." + _))
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/ConsumerProgressTracking.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/ConsumerProgressTracking.scala
index 591ef3fc..c2b90de7 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/ConsumerProgressTracking.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/ConsumerProgressTracking.scala
@@ -51,7 +51,7 @@ trait ConsumerProgressTracking extends
ConsumerAssignmentTrackingListener {
def received[K, V](records: ConsumerRecords[K, V]): Unit = {}
def committed(offsets: java.util.Map[TopicPartition, OffsetAndMetadata]):
Unit = {}
def assignedPositionsAndSeek(assignedTps: Set[TopicPartition],
- consumer: Consumer[_, _],
+ consumer: Consumer[?, ?],
positionTimeout: java.time.Duration): Unit = {}
def addProgressTrackingCallback(callback:
ConsumerAssignmentTrackingListener): Unit = {}
}
@@ -140,7 +140,7 @@ final class ConsumerProgressTrackerImpl extends
ConsumerProgressTracking {
}
override def assignedPositionsAndSeek(assignedTps: Set[TopicPartition],
- consumer: Consumer[_, _],
+ consumer: Consumer[?, ?],
positionTimeout: java.time.Duration): Unit = {
val assignedOffsets = assignedTps.map(tp => tp -> consumer.position(tp,
positionTimeout)).toMap
assignedPositions(assignedTps, assignedOffsets)
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/ControlImplementations.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/ControlImplementations.scala
index 2b192559..3eeefa01 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/ControlImplementations.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/ControlImplementations.scala
@@ -41,7 +41,7 @@ private object PromiseControl {
private trait PromiseControl extends GraphStageLogic with
scaladsl.Consumer.Control {
import PromiseControl._
- def shape: SourceShape[_]
+ def shape: SourceShape[?]
def performShutdown(): Unit
def performStop(): Unit = {
setKeepGoing(true)
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/DefaultProducerStage.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/DefaultProducerStage.scala
index 38912d69..c47519d4 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/DefaultProducerStage.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/DefaultProducerStage.scala
@@ -61,7 +61,7 @@ private class DefaultProducerStageLogic[K, V, P, IN <:
Envelope[K, V, P], OUT <:
override protected def getExecutionContext(): ExecutionContext =
materializer.executionContext
- override protected def logSource: Class[_] = classOf[DefaultProducerStage[_,
_, _, _, _]]
+ override protected def logSource: Class[?] = classOf[DefaultProducerStage[?,
?, ?, ?, ?]]
final override val producerSettings: ProducerSettings[K, V] = stage.settings
@@ -180,7 +180,7 @@ private class DefaultProducerStageLogic[K, V, P, IN <:
Envelope[K, V, P], OUT <:
}
- private abstract class CallbackBase(promise: Promise[_]) extends Callback {
+ private abstract class CallbackBase(promise: Promise[?]) extends Callback {
protected def emitElement(metadata: RecordMetadata): Unit
override def onCompletion(metadata: RecordMetadata, exception: Exception):
Unit =
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/ExternalSingleSourceLogic.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/ExternalSingleSourceLogic.scala
index c3d4630f..29227c69 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/ExternalSingleSourceLogic.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/ExternalSingleSourceLogic.scala
@@ -32,7 +32,7 @@ import scala.concurrent.Future
_consumerActor: ActorRef,
val subscription: ManualSubscription) extends BaseSingleSourceLogic[K, V,
Msg](shape) {
- final override protected def logSource: Class[_] =
classOf[ExternalSingleSourceLogic[K, V, Msg]]
+ final override protected def logSource: Class[?] =
classOf[ExternalSingleSourceLogic[K, V, Msg]]
final val consumerFuture: Future[ActorRef] =
Future.successful(_consumerActor)
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
index 05123265..5f999c16 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
@@ -263,7 +263,7 @@ import scala.util.control.NonFatal
case NonFatal(e) => sendFailure(e, sender())
}
- case p: Poll[_, _] =>
+ case p: Poll[?, ?] =>
receivePoll(p)
case req: RequestMessages =>
@@ -294,7 +294,7 @@ import scala.util.control.NonFatal
case RequestMetrics =>
try {
- val unmodifiableYetMutableMetrics: java.util.Map[MetricName, _ <:
Metric] = consumer.metrics()
+ val unmodifiableYetMutableMetrics: java.util.Map[MetricName, ? <:
Metric] = consumer.metrics()
sender() ! ConsumerMetrics(unmodifiableYetMutableMetrics.asScala.toMap)
} catch {
case NonFatal(e) => sendFailure(e, sender())
@@ -391,7 +391,7 @@ import scala.util.control.NonFatal
}
def stopping: Receive = LoggingReceive.withLabel("stopping") {
- case p: Poll[_, _] =>
+ case p: Poll[?, ?] =>
receivePoll(p)
case _: StopLike =>
case Terminated(ref) =>
@@ -484,7 +484,7 @@ import scala.util.control.NonFatal
commitAndPoll()
}
- private def receivePoll(p: Poll[_, _]): Unit =
+ private def receivePoll(p: Poll[?, ?]): Unit =
if (p.target == this) {
commitAndPoll()
if (p.periodic)
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/MessageBuilder.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/MessageBuilder.scala
index 099f4c04..98c87246 100644
--- a/core/src/main/scala/org/apache/pekko/kafka/internal/MessageBuilder.scala
+++ b/core/src/main/scala/org/apache/pekko/kafka/internal/MessageBuilder.scala
@@ -118,7 +118,7 @@ private[kafka] trait CommittableMessageBuilder[K, V]
extends MessageBuilder[K, V
}
private[kafka] object CommittableMessageBuilder {
- val NoMetadataFromRecord: ConsumerRecord[_, _] => String = (_:
ConsumerRecord[_, _]) =>
+ val NoMetadataFromRecord: ConsumerRecord[?, ?] => String = (_:
ConsumerRecord[?, ?]) =>
OffsetFetchResponse.NO_METADATA
}
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/SingleSourceLogic.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/SingleSourceLogic.scala
index b8cef780..e2b832c6 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/SingleSourceLogic.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/SingleSourceLogic.scala
@@ -35,7 +35,7 @@ import scala.concurrent.{ Future, Promise }
settings: ConsumerSettings[K, V],
override protected val subscription: Subscription) extends
BaseSingleSourceLogic[K, V, Msg](shape) {
- override protected def logSource: Class[_] = classOf[SingleSourceLogic[K, V,
Msg]]
+ override protected def logSource: Class[?] = classOf[SingleSourceLogic[K, V,
Msg]]
private val consumerPromise = Promise[ActorRef]()
final val actorNumber = KafkaConsumerActor.Internal.nextNumber()
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalProducerStage.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalProducerStage.scala
index 435717e0..6766beba 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalProducerStage.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalProducerStage.scala
@@ -124,7 +124,7 @@ private final class TransactionalProducerStageLogic[K, V,
P](
private var firstMessage: Option[Envelope[K, V, P]] = None
- override protected def logSource: Class[_] =
classOf[TransactionalProducerStage[_, _, _]]
+ override protected def logSource: Class[?] =
classOf[TransactionalProducerStage[?, ?, ?]]
// we need to peek at the first message to generate the producer
transactional id for partitioned sources
override def preStart(): Unit = resumeDemand()
diff --git
a/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalSources.scala
b/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalSources.scala
index e712ff98..22e5c660 100644
---
a/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalSources.scala
+++
b/core/src/main/scala/org/apache/pekko/kafka/internal/TransactionalSources.scala
@@ -98,7 +98,7 @@ private[internal] abstract class TransactionalSourceLogic[K,
V, Msg](shape: Sour
import TransactionalSourceLogic._
- override protected def logSource: Class[_] =
classOf[TransactionalSourceLogic[_, _, _]]
+ override protected def logSource: Class[?] =
classOf[TransactionalSourceLogic[?, ?, ?]]
private val inFlightRecords = InFlightRecords.empty
diff --git a/core/src/main/scala/org/apache/pekko/kafka/javadsl/Producer.scala
b/core/src/main/scala/org/apache/pekko/kafka/javadsl/Producer.scala
index 6eacd2a4..a1784f15 100644
--- a/core/src/main/scala/org/apache/pekko/kafka/javadsl/Producer.scala
+++ b/core/src/main/scala/org/apache/pekko/kafka/javadsl/Producer.scala
@@ -84,7 +84,7 @@ object Producer {
* committing, so it is "at-least once delivery" semantics.
*/
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/880")
- def committableSinkWithOffsetContext[K, V, IN <: Envelope[K, V, _], C <:
Committable](
+ def committableSinkWithOffsetContext[K, V, IN <: Envelope[K, V, ?], C <:
Committable](
producerSettings: ProducerSettings[K, V],
committerSettings: CommitterSettings): Sink[pekko.japi.Pair[IN, C],
CompletionStage[Done]] =
committableSink(producerSettings, committerSettings)
diff --git a/core/src/main/scala/org/apache/pekko/kafka/scaladsl/Producer.scala
b/core/src/main/scala/org/apache/pekko/kafka/scaladsl/Producer.scala
index 181f48cf..ee50ef06 100644
--- a/core/src/main/scala/org/apache/pekko/kafka/scaladsl/Producer.scala
+++ b/core/src/main/scala/org/apache/pekko/kafka/scaladsl/Producer.scala
@@ -82,7 +82,7 @@ object Producer {
@ApiMayChange(issue = "https://github.com/akka/alpakka-kafka/issues/880")
def committableSinkWithOffsetContext[K, V](
producerSettings: ProducerSettings[K, V],
- committerSettings: CommitterSettings): Sink[(Envelope[K, V, _],
Committable), Future[Done]] =
+ committerSettings: CommitterSettings): Sink[(Envelope[K, V, ?],
Committable), Future[Done]] =
committableSink(producerSettings, committerSettings)
.contramap {
case (env, offset) =>
diff --git
a/int-tests/src/test/scala/org/apache/pekko/kafka/IntegrationTests.scala
b/int-tests/src/test/scala/org/apache/pekko/kafka/IntegrationTests.scala
index 17ede755..bf3b8b09 100644
--- a/int-tests/src/test/scala/org/apache/pekko/kafka/IntegrationTests.scala
+++ b/int-tests/src/test/scala/org/apache/pekko/kafka/IntegrationTests.scala
@@ -41,8 +41,8 @@ object IntegrationTests {
i
}
- def stopRandomBroker(brokers: Vector[GenericContainer[_]], msgCount:
Long)(implicit log: Logger): Unit = {
- val broker: GenericContainer[_] =
brokers(scala.util.Random.nextInt(brokers.length))
+ def stopRandomBroker(brokers: Vector[GenericContainer[?]], msgCount:
Long)(implicit log: Logger): Unit = {
+ val broker: GenericContainer[?] =
brokers(scala.util.Random.nextInt(brokers.length))
val id = broker.getContainerId
val networkAliases = broker.getNetworkAliases.asScala.mkString(",")
log.warn(
diff --git a/project/CopyrightHeader.scala b/project/CopyrightHeader.scala
index 082caff9..25a34f20 100644
--- a/project/CopyrightHeader.scala
+++ b/project/CopyrightHeader.scala
@@ -23,7 +23,7 @@ trait CopyrightHeader extends AutoPlugin {
override def trigger: PluginTrigger = allRequirements
- protected def headerMappingSettings: Seq[Def.Setting[_]] =
+ protected def headerMappingSettings: Seq[Def.Setting[?]] =
Seq(Compile, Test, Default).flatMap { config =>
inConfig(config)(
Seq(
@@ -35,9 +35,9 @@ trait CopyrightHeader extends AutoPlugin {
HeaderFileType("template") -> cStyleComment)))
}
- override def projectSettings: Seq[Def.Setting[_]] =
Def.settings(headerMappingSettings, additional)
+ override def projectSettings: Seq[Def.Setting[?]] =
Def.settings(headerMappingSettings, additional)
- def additional: Seq[Def.Setting[_]] =
+ def additional: Seq[Def.Setting[?]] =
Def.settings(Compile / compile := {
(Compile / headerCreate).value
(Compile / compile).value
diff --git a/project/CopyrightHeaderForBuild.scala
b/project/CopyrightHeaderForBuild.scala
index d26c4fdc..54c7d048 100644
--- a/project/CopyrightHeaderForBuild.scala
+++ b/project/CopyrightHeaderForBuild.scala
@@ -18,7 +18,7 @@ import sbt.{ inConfig, Compile, Def, PluginTrigger, Test, _ }
object CopyrightHeaderForBuild extends CopyrightHeader {
override def trigger: PluginTrigger = noTrigger
- override def projectSettings: Seq[Def.Setting[_]] = {
+ override def projectSettings: Seq[Def.Setting[?]] = {
Seq(Compile, Test).flatMap { config =>
inConfig(config) {
Seq(
diff --git a/project/CopyrightHeaderForProtobuf.scala
b/project/CopyrightHeaderForProtobuf.scala
index b5d63e5a..a7096054 100644
--- a/project/CopyrightHeaderForProtobuf.scala
+++ b/project/CopyrightHeaderForProtobuf.scala
@@ -16,7 +16,7 @@ import sbt.Keys.sourceDirectory
import sbt.{ inConfig, Compile, Def, Test, _ }
object CopyrightHeaderForProtobuf extends CopyrightHeader {
- override protected def headerMappingSettings: Seq[Def.Setting[_]] = {
+ override protected def headerMappingSettings: Seq[Def.Setting[?]] = {
super.headerMappingSettings
Seq(Compile, Test).flatMap { config =>
inConfig(config) {
diff --git a/project/ProjectSettings.scala b/project/ProjectSettings.scala
index 1e9aed2e..32b7b34a 100644
--- a/project/ProjectSettings.scala
+++ b/project/ProjectSettings.scala
@@ -66,7 +66,7 @@ object ProjectSettings extends AutoPlugin {
private val apacheBaseRepo = "repository.apache.org"
- lazy val commonSettings: Seq[Def.Setting[_]] = Def.settings(
+ lazy val commonSettings: Seq[Def.Setting[?]] = Def.settings(
homepage :=
Some(url("https://pekko.apache.org/docs/pekko-connectors-kafka/current/")),
scmInfo :=
Some(ScmInfo(url("https://github.com/apache/pekko-connectors-kafka"),
"[email protected]:apache/pekko-connectors-kafka.git")),
diff --git a/project/VersionGenerator.scala b/project/VersionGenerator.scala
index c3ba9455..1afb5b6c 100644
--- a/project/VersionGenerator.scala
+++ b/project/VersionGenerator.scala
@@ -17,7 +17,7 @@ import sbt.Keys._
*/
object VersionGenerator {
- val settings: Seq[Setting[_]] = inConfig(Compile)(
+ val settings: Seq[Setting[?]] = inConfig(Compile)(
Seq(
resourceGenerators += generateVersion(resourceManaged, _ /
"version.conf", """|pekko.kafka.version = "%s"
|"""),
diff --git
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/KafkaTestkitTestcontainersSettings.scala
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/KafkaTestkitTestcontainersSettings.scala
index 94cc019b..2367e43e 100644
---
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/KafkaTestkitTestcontainersSettings.scala
+++
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/KafkaTestkitTestcontainersSettings.scala
@@ -44,12 +44,12 @@ final class KafkaTestkitTestcontainersSettings private (
new Consumer[java.util.Collection[PekkoConnectorsKafkaContainer]]() {
override def accept(arg:
java.util.Collection[PekkoConnectorsKafkaContainer]): Unit = ()
},
- val configureZooKeeper: GenericContainer[_] => Unit = _ => (),
- val configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[_]] =
- new Consumer[GenericContainer[_]]() {
- override def accept(arg: GenericContainer[_]): Unit = ()
+ val configureZooKeeper: GenericContainer[?] => Unit = _ => (),
+ val configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[?]] =
+ new Consumer[GenericContainer[?]]() {
+ override def accept(arg: GenericContainer[?]): Unit = ()
},
- val configureSchemaRegistry: GenericContainer[_] => Unit = _ => ()) {
+ val configureSchemaRegistry: GenericContainer[?] => Unit = _ => ()) {
/**
* Java Api
@@ -178,7 +178,7 @@ final class KafkaTestkitTestcontainersSettings private (
/**
* Replaces the default ZooKeeper testcontainers configuration logic
*/
- def withConfigureZooKeeper(configureZooKeeper: GenericContainer[_] => Unit):
KafkaTestkitTestcontainersSettings =
+ def withConfigureZooKeeper(configureZooKeeper: GenericContainer[?] => Unit):
KafkaTestkitTestcontainersSettings =
copy(configureZooKeeper = configureZooKeeper)
/**
@@ -187,7 +187,7 @@ final class KafkaTestkitTestcontainersSettings private (
* Replaces the default ZooKeeper testcontainers configuration logic
*/
def withConfigureZooKeeperConsumer(
- configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[_]])
+ configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[?]])
: KafkaTestkitTestcontainersSettings =
copy(configureZooKeeperConsumer = configureZooKeeperConsumer)
@@ -196,7 +196,7 @@ final class KafkaTestkitTestcontainersSettings private (
* Replaces the default schema registry testcontainers configuration logic
*/
def withConfigureSchemaRegistry(
- configureSchemaRegistry: GenericContainer[_] => Unit):
KafkaTestkitTestcontainersSettings =
+ configureSchemaRegistry: GenericContainer[?] => Unit):
KafkaTestkitTestcontainersSettings =
copy(configureSchemaRegistry = configureSchemaRegistry)
/**
@@ -255,9 +255,9 @@ final class KafkaTestkitTestcontainersSettings private (
configureKafka: Vector[PekkoConnectorsKafkaContainer] => Unit =
configureKafka,
configureKafkaConsumer:
java.util.function.Consumer[java.util.Collection[PekkoConnectorsKafkaContainer]]
=
configureKafkaConsumer,
- configureZooKeeper: GenericContainer[_] => Unit = configureZooKeeper,
- configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[_]] = configureZooKeeperConsumer,
- configureSchemaRegistry: GenericContainer[_] => Unit =
configureSchemaRegistry)
+ configureZooKeeper: GenericContainer[?] => Unit = configureZooKeeper,
+ configureZooKeeperConsumer:
java.util.function.Consumer[GenericContainer[?]] = configureZooKeeperConsumer,
+ configureSchemaRegistry: GenericContainer[?] => Unit =
configureSchemaRegistry)
: KafkaTestkitTestcontainersSettings =
new KafkaTestkitTestcontainersSettings(zooKeeperImage,
zooKeeperImageTag,
diff --git
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/ProducerResultFactory.scala
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/ProducerResultFactory.scala
index 66571808..12986189 100644
---
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/ProducerResultFactory.scala
+++
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/ProducerResultFactory.scala
@@ -29,7 +29,7 @@ import scala.jdk.CollectionConverters._
@ApiMayChange
object ProducerResultFactory {
- def recordMetadata(msg: ProducerRecord[_, _]): RecordMetadata = {
+ def recordMetadata(msg: ProducerRecord[?, ?]): RecordMetadata = {
// null checks are required on Scala 2.11
val partition = if (msg.partition == null) 0 else msg.partition.toInt
val timestamp = if (msg.timestamp == null) 0L else msg.timestamp.toLong
diff --git
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/internal/TestcontainersKafka.scala
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/internal/TestcontainersKafka.scala
index 32dc338f..4d2c6e82 100644
---
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/internal/TestcontainersKafka.scala
+++
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/internal/TestcontainersKafka.scala
@@ -54,7 +54,7 @@ object TestcontainersKafka {
cluster.getBrokers.asScala.toVector
}
- def zookeeperContainer: Option[GenericContainer[_]] = {
+ def zookeeperContainer: Option[GenericContainer[?]] = {
requireStarted()
cluster.getZooKeeper.toScala
}
diff --git
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/scaladsl/TestcontainersKafkaLike.scala
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/scaladsl/TestcontainersKafkaLike.scala
index 1b0a8f60..758a3dac 100644
---
a/testkit/src/main/scala/org/apache/pekko/kafka/testkit/scaladsl/TestcontainersKafkaLike.scala
+++
b/testkit/src/main/scala/org/apache/pekko/kafka/testkit/scaladsl/TestcontainersKafkaLike.scala
@@ -30,7 +30,7 @@ trait TestcontainersKafkaLike extends
TestcontainersKafka.Spec {
override def kafkaPort: Int = TestcontainersKafka.Singleton.kafkaPort
override def bootstrapServers: String =
TestcontainersKafka.Singleton.bootstrapServers
override def brokerContainers: Vector[PekkoConnectorsKafkaContainer] =
TestcontainersKafka.Singleton.brokerContainers
- override def zookeeperContainer: Option[GenericContainer[_]] =
TestcontainersKafka.Singleton.zookeeperContainer
+ override def zookeeperContainer: Option[GenericContainer[?]] =
TestcontainersKafka.Singleton.zookeeperContainer
override def schemaRegistryContainer: Option[SchemaRegistryContainer] =
TestcontainersKafka.Singleton.schemaRegistryContainer
override def schemaRegistryUrl: String =
TestcontainersKafka.Singleton.schemaRegistryUrl
diff --git a/tests/src/test/scala/docs/scaladsl/ClusterShardingExample.scala
b/tests/src/test/scala/docs/scaladsl/ClusterShardingExample.scala
index 0126bb06..70443d74 100644
--- a/tests/src/test/scala/docs/scaladsl/ClusterShardingExample.scala
+++ b/tests/src/test/scala/docs/scaladsl/ClusterShardingExample.scala
@@ -39,7 +39,7 @@ import scala.util.{ Failure, Success }
*
https://github.com/akka/akka-samples/tree/2.6/akka-sample-kafka-to-sharding-scala
*/
object ClusterShardingExample {
- implicit val system: ActorSystem[_] = ActorSystem(Behaviors.empty,
"ClusterShardingExample")
+ implicit val system: ActorSystem[?] = ActorSystem(Behaviors.empty,
"ClusterShardingExample")
val kafkaBootstrapServers = "localhost:9092"
implicit val ec: ExecutionContext = system.executionContext
diff --git a/tests/src/test/scala/docs/scaladsl/ProducerExample.scala
b/tests/src/test/scala/docs/scaladsl/ProducerExample.scala
index cb81e158..a068a0e2 100644
--- a/tests/src/test/scala/docs/scaladsl/ProducerExample.scala
+++ b/tests/src/test/scala/docs/scaladsl/ProducerExample.scala
@@ -113,7 +113,7 @@ class ProducerExample extends DocsSpecBase with
TestcontainersKafkaLike {
.createKafkaProducerAsync()
.map { kafkaProducer =>
// #producerMetrics
- val metrics: java.util.Map[org.apache.kafka.common.MetricName, _ <:
org.apache.kafka.common.Metric] =
+ val metrics: java.util.Map[org.apache.kafka.common.MetricName, ? <:
org.apache.kafka.common.Metric] =
kafkaProducer.metrics() // observe metrics
// #producerMetrics
metrics.isEmpty should be(false)
diff --git a/tests/src/test/scala/docs/scaladsl/proto/Order.scala
b/tests/src/test/scala/docs/scaladsl/proto/Order.scala
index c3d9495e..dfccb3ac 100644
--- a/tests/src/test/scala/docs/scaladsl/proto/Order.scala
+++ b/tests/src/test/scala/docs/scaladsl/proto/Order.scala
@@ -120,11 +120,11 @@ object Order extends
scalapb.GeneratedMessageCompanion[docs.scaladsl.proto.Order
def javaDescriptor: _root_.com.google.protobuf.Descriptors.Descriptor =
OrderProto.javaDescriptor.getMessageTypes.get(0)
def scalaDescriptor: _root_.scalapb.descriptors.Descriptor =
OrderProto.scalaDescriptor.messages(0)
- def messageCompanionForFieldNumber(__number: _root_.scala.Int):
_root_.scalapb.GeneratedMessageCompanion[_] =
+ def messageCompanionForFieldNumber(__number: _root_.scala.Int):
_root_.scalapb.GeneratedMessageCompanion[?] =
throw new MatchError(__number)
lazy val nestedMessagesCompanions
- : Seq[_root_.scalapb.GeneratedMessageCompanion[_ <:
_root_.scalapb.GeneratedMessage]] = Seq.empty
- def enumCompanionForFieldNumber(__fieldNumber: _root_.scala.Int):
_root_.scalapb.GeneratedEnumCompanion[_] =
+ : Seq[_root_.scalapb.GeneratedMessageCompanion[? <:
_root_.scalapb.GeneratedMessage]] = Seq.empty
+ def enumCompanionForFieldNumber(__fieldNumber: _root_.scala.Int):
_root_.scalapb.GeneratedEnumCompanion[?] =
throw new MatchError(__fieldNumber)
lazy val defaultInstance = docs.scaladsl.proto.Order(
id = "")
diff --git a/tests/src/test/scala/docs/scaladsl/proto/OrderProto.scala
b/tests/src/test/scala/docs/scaladsl/proto/OrderProto.scala
index b26df6b9..829f93e8 100644
--- a/tests/src/test/scala/docs/scaladsl/proto/OrderProto.scala
+++ b/tests/src/test/scala/docs/scaladsl/proto/OrderProto.scala
@@ -21,8 +21,8 @@ package docs.scaladsl.proto
object OrderProto extends _root_.scalapb.GeneratedFileObject {
lazy val dependencies: Seq[_root_.scalapb.GeneratedFileObject] = Seq.empty
- lazy val messagesCompanions: Seq[_root_.scalapb.GeneratedMessageCompanion[_
<: _root_.scalapb.GeneratedMessage]] =
- Seq[_root_.scalapb.GeneratedMessageCompanion[_ <:
_root_.scalapb.GeneratedMessage]](
+ lazy val messagesCompanions: Seq[_root_.scalapb.GeneratedMessageCompanion[?
<: _root_.scalapb.GeneratedMessage]] =
+ Seq[_root_.scalapb.GeneratedMessageCompanion[? <:
_root_.scalapb.GeneratedMessage]](
docs.scaladsl.proto.Order)
private lazy val ProtoBytes: Array[Byte] =
scalapb.Encoding.fromBase64(
diff --git
a/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerDummy.scala
b/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerDummy.scala
index aa96dfb3..9dbb1ce5 100644
--- a/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerDummy.scala
+++ b/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerDummy.scala
@@ -63,7 +63,7 @@ abstract class ConsumerDummy[K, V] extends Consumer[K, V] {
override def seekToEnd(partitions: java.util.Collection[TopicPartition]):
Unit = ???
override def position(partition: TopicPartition): Long = ???
override def position(partition: TopicPartition, timeout:
java.time.Duration): Long = ???
- override def metrics(): java.util.Map[MetricName, _ <: Metric] = ???
+ override def metrics(): java.util.Map[MetricName, ? <: Metric] = ???
override def partitionsFor(topic: String): java.util.List[PartitionInfo] =
???
override def listTopics(): java.util.Map[String,
java.util.List[PartitionInfo]] = ???
override def paused(): java.util.Set[TopicPartition] = ???
diff --git
a/tests/src/test/scala/org/apache/pekko/kafka/scaladsl/ReconnectSpec.scala
b/tests/src/test/scala/org/apache/pekko/kafka/scaladsl/ReconnectSpec.scala
index 22fa4321..9b0300fb 100644
--- a/tests/src/test/scala/org/apache/pekko/kafka/scaladsl/ReconnectSpec.scala
+++ b/tests/src/test/scala/org/apache/pekko/kafka/scaladsl/ReconnectSpec.scala
@@ -48,7 +48,7 @@ class ReconnectSpec extends SpecBase with
TestcontainersKafkaLike {
.to(Producer.plainSink(producerDefaults.withBootstrapServers(s"localhost:$proxyPort")))
.run()
- def offerInOrder(msgs: Seq[String]): Future[_] =
+ def offerInOrder(msgs: Seq[String]): Future[?] =
if (msgs.isEmpty) Future.successful(Done)
else producer.offer(msgs.head).flatMap(_ => offerInOrder(msgs.tail))
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]