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.git


The following commit(s) were added to refs/heads/main by this push:
     new da1a39e3c fix: Scala 3.8+ forward compatibility (#1682)
da1a39e3c is described below

commit da1a39e3cee59cbc1eb15e2d4ee6881acd843976
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Mon Jun 15 17:12:14 2026 +0800

    fix: Scala 3.8+ forward compatibility (#1682)
    
    * fix: Scala 3.8+ forward compatibility
    
    Motivation:
    Prepare for Scala 3.8/3.9 adoption by fixing deprecation warnings
    that appear with newer Scala compilers while maintaining compatibility
    with current Scala 3.3.x and 2.13.x.
    
    Modification:
    - Replace `= _` with explicit `= null`/`= 0`/`= null.asInstanceOf[T]`
      in var declarations (deprecated in 3.8)
    - Replace `private[this]` with `private` (deprecated in 3.8)
    - Change context bounds to explicit implicit parameters in Google
      connector methods that are called with explicit args from Java API:
      GoogleHttp.cachedHostConnectionPool,
      GoogleHttp.cachedHostConnectionPoolWithContext,
      scaladsl.Google.singleRequest, scaladsl.Google.resumableUpload
    - Move Scala 2-only compiler options (-Xlint, -Ywarn-dead-code,
      cat=unused-nowarn) to Scala 2 conditional block
    - Add -Wconf suppressions for cross-version incompatible warnings
    
    Result:
    Clean compilation on Scala 3.8.4 with zero warnings.
    No Scala version changes in production build.
    
    * fix: add MiMa exclusion for BigQuery.singleRequest context bounds change
    
    Motivation:
    Converting context bound [T: FromResponseUnmarshaller] to explicit
    implicit parameter reorders the implicit parameter list at the
    bytecode level, causing MiMa to report IncompatibleMethTypeProblem.
    
    Modification:
    Add MiMa exclusion filter for BigQuery.singleRequest. The method is
    @ApiMayChange and source-level callers using implicit resolution are
    unaffected.
    
    Result:
    MiMa binary compatibility check passes.
    
    Tests:
    Not run - MiMa filter addition only
    
    References:
    Fixes #1682
---
 .../connectors/amqp/impl/AmqpConnectorLogic.scala  |  4 +--
 .../connectors/amqp/impl/AmqpRpcFlowStage.scala    |  2 +-
 .../aws/eventbridge/IntegrationTestContext.scala   |  4 +--
 .../pekko/stream/connectors/couchbase/model.scala  |  2 +-
 .../couchbase/testing/CouchbaseSupport.scala       |  2 +-
 .../stream/connectors/csv/scaladsl/CsvBench.scala  |  4 +--
 .../stream/connectors/csv/impl/CsvFormatter.scala  | 12 +++----
 .../stream/connectors/csv/impl/CsvParser.scala     | 42 +++++++++++-----------
 .../connectors/csv/impl/CsvParsingStage.scala      |  2 +-
 .../connectors/csv/impl/CsvToMapJavaStage.scala    |  2 +-
 .../connectors/csv/scaladsl/ByteOrderMark.scala    |  2 +-
 .../connectors/file/scaladsl/LogRotatorSink.scala  |  2 +-
 .../connectors/ftp/impl/FtpBrowserGraphStage.scala | 10 +++---
 .../connectors/ftp/impl/FtpIOGraphStage.scala      | 12 +++----
 .../connectors/ftp/impl/SftpOperations.scala       |  2 +-
 .../geode/impl/stage/GeodeCQueryGraphLogic.scala   |  2 +-
 .../geode/impl/stage/GeodeSourceStageLogic.scala   |  2 +-
 .../bigquery/storage/javadsl/BigQueryStorage.scala |  4 +--
 .../context-bounds.excludes                        | 22 ++++++++++++
 .../impl/GCStorageStreamIntegrationSpec.scala      |  2 +-
 .../stream/connectors/google/http/GoogleHttp.scala | 10 +++---
 .../stream/connectors/google/javadsl/Google.scala  |  2 +-
 .../connectors/google/jwt/JwtSprayJson.scala       |  4 +--
 .../stream/connectors/google/scaladsl/Google.scala |  8 +++--
 .../test/scala/docs/scaladsl/HdfsReaderSpec.scala  |  2 +-
 .../test/scala/docs/scaladsl/HdfsWriterSpec.scala  |  2 +-
 .../influxdb/impl/InfluxDbSourceStage.scala        |  2 +-
 .../src/test/scala/docs/scaladsl/FlowSpec.scala    |  2 +-
 .../scala/docs/scaladsl/InfluxDbSourceSpec.scala   |  2 +-
 .../test/scala/docs/scaladsl/InfluxDbSpec.scala    |  2 +-
 .../connectors/ironmq/impl/IronMqPullStage.scala   |  2 +-
 .../connectors/ironmq/impl/IronMqPushStage.scala   |  2 +-
 .../stream/connectors/jakartams/Envelopes.scala    |  2 +-
 .../connectors/jakartams/impl/JmsBrowseStage.scala |  8 ++---
 .../connectors/jakartams/impl/JmsConnector.scala   |  4 +--
 .../connectors/jakartams/JmsSharedServerSpec.scala |  4 +--
 .../pekko/stream/connectors/jms/Envelopes.scala    |  2 +-
 .../connectors/jms/impl/JmsBrowseStage.scala       |  8 ++---
 .../stream/connectors/jms/impl/JmsConnector.scala  |  4 +--
 .../connectors/jms/JmsSharedServerSpec.scala       |  4 +--
 .../kinesis/impl/KinesisSchedulerSourceStage.scala |  6 ++--
 .../kinesis/impl/KinesisSourceStage.scala          | 12 +++----
 .../connectors/kinesis/impl/ShardProcessor.scala   |  4 +--
 .../kinesis/KinesisSchedulerSourceSpec.scala       | 10 +++---
 .../orientdb/impl/OrientDbFlowStage.scala          |  4 +--
 .../orientdb/impl/OrientDbSourceStage.scala        |  4 +--
 .../test/scala/docs/scaladsl/OrientDbSpec.scala    |  6 ++--
 .../connectors/pravega/impl/PravegaFlow.scala      |  2 +-
 .../connectors/pravega/impl/PravegaSource.scala    |  2 +-
 .../pravega/impl/PravegaTableReadFlow.scala        |  4 +--
 .../pravega/impl/PravegaTableSource.scala          |  4 +--
 .../pravega/impl/PravegaTableWriteFlow.scala       |  4 +--
 project/Common.scala                               | 25 +++++++++----
 .../stream/connectors/s3/impl/HttpRequests.scala   |  6 ++--
 .../stream/connectors/s3/impl/auth/Signer.scala    |  6 ++--
 .../connectors/sns/IntegrationTestContext.scala    |  4 +--
 solr/src/test/scala/docs/scaladsl/SolrSpec.scala   |  4 +--
 .../connectors/sqs/impl/BalancingMapAsync.scala    |  2 +-
 .../pekko/stream/connectors/udp/impl/UdpBind.scala |  2 +-
 .../pekko/stream/connectors/udp/impl/UdpSend.scala |  2 +-
 .../connectors/xml/impl/StreamingXmlParser.scala   |  2 +-
 61 files changed, 183 insertions(+), 144 deletions(-)

diff --git 
a/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpConnectorLogic.scala
 
b/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpConnectorLogic.scala
index e50fa7736..cf237b0ca 100644
--- 
a/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpConnectorLogic.scala
+++ 
b/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpConnectorLogic.scala
@@ -22,8 +22,8 @@ import scala.util.control.NonFatal
 
 private trait AmqpConnectorLogic { this: GraphStageLogic =>
 
-  private var connection: Connection = _
-  protected var channel: Channel = _
+  private var connection: Connection = null
+  protected var channel: Channel = null
 
   protected lazy val shutdownCallback: AsyncCallback[Throwable] = 
getAsyncCallback(onFailure)
   private lazy val shutdownListener = new ShutdownListener {
diff --git 
a/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpRpcFlowStage.scala
 
b/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpRpcFlowStage.scala
index e2b186336..3ca4127fd 100644
--- 
a/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpRpcFlowStage.scala
+++ 
b/amqp/src/main/scala/org/apache/pekko/stream/connectors/amqp/impl/AmqpRpcFlowStage.scala
@@ -58,7 +58,7 @@ private[amqp] final class AmqpRpcFlowStage(writeSettings: 
AmqpWriteSettings, buf
         private val exchange = settings.exchange.getOrElse("")
         private val routingKey = settings.routingKey.getOrElse("")
         private val queue = mutable.Queue[CommittableReadResult]()
-        private var queueName: String = _
+        private var queueName: String = null
         private var unackedMessages = 0
         private var outstandingMessages = 0
 
diff --git 
a/aws-event-bridge/src/test/scala/org/apache/pekko/stream/connectors/aws/eventbridge/IntegrationTestContext.scala
 
b/aws-event-bridge/src/test/scala/org/apache/pekko/stream/connectors/aws/eventbridge/IntegrationTestContext.scala
index 338554ea2..cd72e2688 100644
--- 
a/aws-event-bridge/src/test/scala/org/apache/pekko/stream/connectors/aws/eventbridge/IntegrationTestContext.scala
+++ 
b/aws-event-bridge/src/test/scala/org/apache/pekko/stream/connectors/aws/eventbridge/IntegrationTestContext.scala
@@ -34,8 +34,8 @@ trait IntegrationTestContext extends BeforeAndAfterAll with 
ScalaFutures {
 
   def eventBusEndpoint: String = s"http://localhost:4587";
 
-  implicit var eventBridgeClient: EventBridgeAsyncClient = _
-  var eventBusArn: String = _
+  implicit var eventBridgeClient: EventBridgeAsyncClient = null
+  var eventBusArn: String = null
 
   def createEventBus(): String =
     eventBridgeClient
diff --git 
a/couchbase/src/main/scala/org/apache/pekko/stream/connectors/couchbase/model.scala
 
b/couchbase/src/main/scala/org/apache/pekko/stream/connectors/couchbase/model.scala
index af17b70f8..e8182a07a 100644
--- 
a/couchbase/src/main/scala/org/apache/pekko/stream/connectors/couchbase/model.scala
+++ 
b/couchbase/src/main/scala/org/apache/pekko/stream/connectors/couchbase/model.scala
@@ -85,7 +85,7 @@ final class CouchbaseWriteSettings private (val parallelism: 
Int,
    */
   def withTimeout(timeout: FiniteDuration): CouchbaseWriteSettings = 
copy(timeout = timeout)
 
-  private[this] def copy(parallelism: Int = parallelism,
+  private def copy(parallelism: Int = parallelism,
       replicateTo: ReplicateTo = replicateTo,
       persistTo: PersistTo = persistTo,
       timeout: FiniteDuration = timeout) =
diff --git 
a/couchbase/src/test/scala/org/apache/pekko/stream/connectors/couchbase/testing/CouchbaseSupport.scala
 
b/couchbase/src/test/scala/org/apache/pekko/stream/connectors/couchbase/testing/CouchbaseSupport.scala
index c780893c4..2aa72d95e 100644
--- 
a/couchbase/src/test/scala/org/apache/pekko/stream/connectors/couchbase/testing/CouchbaseSupport.scala
+++ 
b/couchbase/src/test/scala/org/apache/pekko/stream/connectors/couchbase/testing/CouchbaseSupport.scala
@@ -63,7 +63,7 @@ trait CouchbaseSupport {
   val bucketName = "pekko"
   val queryBucketName = "pekkoquery"
 
-  var session: CouchbaseSession = _
+  var session: CouchbaseSession = null
 
   def beforeAll(): Unit = {
     session = Await.result(CouchbaseSession(sessionSettings, bucketName), 
10.seconds)
diff --git 
a/csv-bench/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/CsvBench.scala
 
b/csv-bench/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/CsvBench.scala
index d81dc4a9e..657ec8b2f 100644
--- 
a/csv-bench/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/CsvBench.scala
+++ 
b/csv-bench/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/CsvBench.scala
@@ -75,8 +75,8 @@ class CsvBench {
       "8192", // ~same size as row
       "65536" // ~8k larger than row
     ))
-  var bsSize: Int = _
-  var source: Source[ByteString, NotUsed] = _
+  var bsSize: Int = 0
+  var source: Source[ByteString, NotUsed] = null
 
   @Benchmark
   def parse(bh: Blackhole): Unit = {
diff --git 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvFormatter.scala
 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvFormatter.scala
index 36f248a74..8d61334d4 100644
--- 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvFormatter.scala
+++ 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvFormatter.scala
@@ -32,13 +32,13 @@ import scala.collection.immutable
     quotingStyle: CsvQuotingStyle,
     charset: Charset = StandardCharsets.UTF_8) {
 
-  private[this] val charsetName = charset.name()
+  private val charsetName = charset.name()
 
-  private[this] val delimiterBs = ByteString(String.valueOf(delimiter), 
charsetName)
-  private[this] val quoteBs = ByteString(String.valueOf(quoteChar), 
charsetName)
-  private[this] val duplicatedQuote = 
ByteString(String.valueOf(Array(quoteChar, quoteChar)), charsetName)
-  private[this] val duplicatedEscape = 
ByteString(String.valueOf(Array(escapeChar, escapeChar)), charsetName)
-  private[this] val endOfLineBs = ByteString(endOfLine, charsetName)
+  private val delimiterBs = ByteString(String.valueOf(delimiter), charsetName)
+  private val quoteBs = ByteString(String.valueOf(quoteChar), charsetName)
+  private val duplicatedQuote = ByteString(String.valueOf(Array(quoteChar, 
quoteChar)), charsetName)
+  private val duplicatedEscape = ByteString(String.valueOf(Array(escapeChar, 
escapeChar)), charsetName)
+  private val endOfLineBs = ByteString(endOfLine, charsetName)
 
   def toCsv(fields: immutable.Iterable[Any]): ByteString =
     if (fields.nonEmpty) nonEmptyToCsv(fields)
diff --git 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParser.scala
 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParser.scala
index 2c229af33..2cd047c30 100644
--- 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParser.scala
+++ 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParser.scala
@@ -60,12 +60,12 @@ import scala.collection.mutable
    *
    * May include previous chunks that start a field but do not complete it.
    */
-  private[this] var buffer: ByteString = ByteString.empty
+  private var buffer: ByteString = ByteString.empty
 
   /**
    * Flag to run BOM checks against first two bytes of the stream.
    */
-  private[this] var firstData = true
+  private var firstData = true
 
   /**
    * Current position within [[buffer]].
@@ -73,7 +73,7 @@ import scala.collection.mutable
    * Points to the same byte as [[current.head]].
    * Used for slicing fields out of [[buffer]] and for debug info.
    */
-  private[this] var pos: Int = 0
+  private var pos: Int = 0
 
   /**
    * Number of bytes dropped on the current row.
@@ -84,27 +84,27 @@ import scala.collection.mutable
    * [[pekko.util.ByteString.ByteStrings]] to a 
[[pekko.util.ByteString.ByteString1]]
    * to exploit the much faster [[ByteString.slice()]] implementation.
    */
-  private[this] var lineBytesDropped = 0
+  private var lineBytesDropped = 0
 
   /**
    * Position within the current row.
    *
    * Used for enforcing line length limits and as debug info for exceptions.
    */
-  private[this] def lineLength: Int = lineBytesDropped + pos
+  private def lineLength: Int = lineBytesDropped + pos
 
   /**
    * Position within [[buffer]] of the start of the current field.
    */
-  private[this] var fieldStart = 0
-  private[this] var currentLineNo = 1L
+  private var fieldStart = 0
+  private var currentLineNo = 1L
 
   /**
    * Reset after each row.
    */
-  private[this] val columns = mutable.ListBuffer[ByteString]()
-  private[this] var state: State = LineStart
-  private[this] val fieldBuilder = new FieldBuilder
+  private val columns = mutable.ListBuffer[ByteString]()
+  private var state: State = LineStart
+  private val fieldBuilder = new FieldBuilder
 
   /**
    * Current iterator being parsed.
@@ -114,7 +114,7 @@ import scala.collection.mutable
    *
    * We fully parse each chunk before getting the next, so we only need to 
track one [[ByteIterator]] at a time.
    */
-  private[this] var current: ByteIterator = ByteString.empty.iterator
+  private var current: ByteIterator = ByteString.empty.iterator
 
   def offer(next: ByteString): Unit =
     if (next.nonEmpty) {
@@ -137,17 +137,17 @@ import scala.collection.mutable
     line
   }
 
-  private[this] def advance(n: Int = 1): Unit = {
+  private def advance(n: Int = 1): Unit = {
     pos += n
     current.drop(n)
   }
 
-  private[this] def resetLine(): Unit = {
+  private def resetLine(): Unit = {
     dropReadBuffer()
     lineBytesDropped = 0
   }
 
-  private[this] def dropReadBuffer() = {
+  private def dropReadBuffer() = {
     buffer = buffer.drop(pos)
     lineBytesDropped += pos
     pos = 0
@@ -163,8 +163,8 @@ import scala.collection.mutable
     /**
      * false if [[builder]] is null.
      */
-    private[this] var useBuilder = false
-    private[this] var builder: ByteStringBuilder = _
+    private var useBuilder = false
+    private var builder: ByteStringBuilder = null
 
     /**
      * Set up the ByteString builder instead of relying on `ByteString.slice`.
@@ -186,12 +186,12 @@ import scala.collection.mutable
 
   }
 
-  private[this] def noCharEscaped() =
+  private def noCharEscaped() =
     throw new MalformedCsvException(currentLineNo,
       lineLength,
       s"wrong escaping at $currentLineNo:$lineLength, no character after 
escape")
 
-  private[this] def checkForByteOrderMark(): Unit =
+  private def checkForByteOrderMark(): Unit =
     if (buffer.length >= 2) {
       if (buffer.startsWith(ByteOrderMark.UTF_8)) {
         advance(3)
@@ -209,7 +209,7 @@ import scala.collection.mutable
       }
     }
 
-  private[this] def parseLine(): Unit = {
+  private def parseLine(): Unit = {
     if (firstData) {
       checkForByteOrderMark()
       firstData = false
@@ -217,7 +217,7 @@ import scala.collection.mutable
     churn()
   }
 
-  private[this] def churn(): Unit = {
+  private def churn(): Unit = {
     while (state != LineEnd && pos < buffer.length) {
       if (lineLength >= maximumLineLength)
         throw new MalformedCsvException(
@@ -405,7 +405,7 @@ import scala.collection.mutable
       }
     }
   }
-  private[this] def maybeExtractLine(requireLineEnd: Boolean): 
Option[List[ByteString]] =
+  private def maybeExtractLine(requireLineEnd: Boolean): 
Option[List[ByteString]] =
     if (requireLineEnd) {
       state match {
         case LineEnd =>
diff --git 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParsingStage.scala
 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParsingStage.scala
index 09f3bd8e0..42aee3cb3 100644
--- 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParsingStage.scala
+++ 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvParsingStage.scala
@@ -40,7 +40,7 @@ import scala.util.control.NonFatal
 
   override def createLogic(inheritedAttributes: Attributes) =
     new GraphStageLogic(shape) with InHandler with OutHandler {
-      private[this] val buffer = new CsvParser(delimiter, quoteChar, 
escapeChar, maximumLineLength)
+      private val buffer = new CsvParser(delimiter, quoteChar, escapeChar, 
maximumLineLength)
 
       setHandlers(in, out, this)
 
diff --git 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvToMapJavaStage.scala
 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvToMapJavaStage.scala
index a7cbd8ca8..87f43a637 100644
--- 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvToMapJavaStage.scala
+++ 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/impl/CsvToMapJavaStage.scala
@@ -57,7 +57,7 @@ import pekko.util.ByteString
 
   override def createLogic(inheritedAttributes: Attributes): GraphStageLogic =
     new GraphStageLogic(shape) {
-      private[this] var headers = columnNames
+      private var headers = columnNames
 
       setHandler(
         in,
diff --git 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/ByteOrderMark.scala
 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/ByteOrderMark.scala
index e270df4a4..d77e2da76 100644
--- 
a/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/ByteOrderMark.scala
+++ 
b/csv/src/main/scala/org/apache/pekko/stream/connectors/csv/scaladsl/ByteOrderMark.scala
@@ -23,7 +23,7 @@ import org.apache.pekko.util.ByteString
  */
 object ByteOrderMark {
 
-  private[this] final val ZeroZero = ByteString.apply(0x00.toByte, 0x00.toByte)
+  private final val ZeroZero = ByteString.apply(0x00.toByte, 0x00.toByte)
 
   /** Byte Order Mark for UTF-16 big-endian */
   final val UTF_16_BE = ByteString.apply(0xFE.toByte, 0xFF.toByte)
diff --git 
a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/LogRotatorSink.scala
 
b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/LogRotatorSink.scala
index b303fa00f..49c972517 100644
--- 
a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/LogRotatorSink.scala
+++ 
b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/LogRotatorSink.scala
@@ -97,7 +97,7 @@ final private class LogRotatorSink[T, C, 
R](triggerGeneratorCreator: () => T =>
 
   private final class Logic(promise: Promise[Done]) extends 
GraphStageLogic(shape) {
     val triggerGenerator: T => Option[C] = triggerGeneratorCreator()
-    var sourceOut: SubSourceOutlet[T] = _
+    var sourceOut: SubSourceOutlet[T] = null
     var sinkCompletions: immutable.Seq[Future[R]] = immutable.Seq.empty
     var isFinishing = false
 
diff --git 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpBrowserGraphStage.scala
 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpBrowserGraphStage.scala
index 589624aa4..74da2c482 100644
--- 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpBrowserGraphStage.scala
+++ 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpBrowserGraphStage.scala
@@ -34,9 +34,9 @@ private[ftp] trait FtpBrowserGraphStage[FtpClient, S <: 
RemoteFileSettings]
   def createLogic(inheritedAttributes: Attributes): 
FtpGraphStageLogic[FtpFile, FtpClient, S] = {
     val logic = new FtpGraphStageLogic[FtpFile, FtpClient, S](shape, ftpLike, 
connectionSettings, ftpClient) {
 
-      private[this] var buffer: Seq[FtpFile] = Seq.empty[FtpFile]
+      private var buffer: Seq[FtpFile] = Seq.empty[FtpFile]
 
-      private[this] var traversed: Seq[FtpFile] = Seq.empty[FtpFile]
+      private var traversed: Seq[FtpFile] = Seq.empty[FtpFile]
 
       setHandler(
         out,
@@ -70,11 +70,11 @@ private[ftp] trait FtpBrowserGraphStage[FtpClient, S <: 
RemoteFileSettings]
 
       override protected[this] def matFailure(t: Throwable) = true
 
-      private[this] def initBuffer(basePath: String) =
+      private def initBuffer(basePath: String) =
         getFilesFromPath(basePath)
 
       @scala.annotation.tailrec
-      private[this] def fillBuffer(): Unit = buffer match {
+      private def fillBuffer(): Unit = buffer match {
         case head +: tail if head.isDirectory && branchSelector(head) =>
           buffer = getFilesFromPath(head.path) ++ tail
           if (emitTraversedDirectories) traversed = traversed :+ head
@@ -82,7 +82,7 @@ private[ftp] trait FtpBrowserGraphStage[FtpClient, S <: 
RemoteFileSettings]
         case _ => // do nothing
       }
 
-      private[this] def getFilesFromPath(basePath: String) =
+      private def getFilesFromPath(basePath: String) =
         if (basePath.isEmpty)
           graphStageFtpLike.listFiles(handler.get)
         else
diff --git 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpIOGraphStage.scala
 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpIOGraphStage.scala
index b26f17426..fa2e346a7 100644
--- 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpIOGraphStage.scala
+++ 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/FtpIOGraphStage.scala
@@ -79,8 +79,8 @@ private[ftp] trait FtpIOSourceStage[FtpClient, S <: 
RemoteFileSettings]
 
     val logic = new FtpGraphStageLogic[ByteString, FtpClient, S](shape, 
ftpLike, connectionSettings, ftpClient) {
 
-      private[this] var isOpt: Option[InputStream] = None
-      private[this] var readBytesTotal: Long = 0L
+      private var isOpt: Option[InputStream] = None
+      private var readBytesTotal: Long = 0L
 
       setHandler(
         out,
@@ -158,7 +158,7 @@ private[ftp] trait FtpIOSourceStage[FtpClient, S <: 
RemoteFileSettings]
         matValuePromise.tryFailure(new 
IOOperationIncompleteException(readBytesTotal, t))
 
       /** BLOCKING I/O READ */
-      private[this] def readChunk() = {
+      private def readChunk() = {
         def read(arr: Array[Byte]) =
           isOpt.flatMap { is =>
             val readBytes = is.read(arr)
@@ -200,8 +200,8 @@ private[ftp] trait FtpIOSinkStage[FtpClient, S <: 
RemoteFileSettings]
 
     val logic = new FtpGraphStageLogic[ByteString, FtpClient, S](shape, 
ftpLike, connectionSettings, ftpClient) {
 
-      private[this] var osOpt: Option[OutputStream] = None
-      private[this] var writtenBytesTotal: Long = 0L
+      private var osOpt: Option[OutputStream] = None
+      private var writtenBytesTotal: Long = 0L
 
       setHandler(
         in,
@@ -262,7 +262,7 @@ private[ftp] trait FtpIOSinkStage[FtpClient, S <: 
RemoteFileSettings]
         matValuePromise.tryFailure(new 
IOOperationIncompleteException(writtenBytesTotal, t))
 
       /** BLOCKING I/O WRITE */
-      private[this] def write(bytes: ByteString) =
+      private def write(bytes: ByteString) =
         osOpt.foreach { os =>
           os.write(bytes.toArray)
           writtenBytesTotal += bytes.size
diff --git 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/SftpOperations.scala
 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/SftpOperations.scala
index e8b75160d..05ee5f538 100644
--- 
a/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/SftpOperations.scala
+++ 
b/ftp/src/main/scala/org/apache/pekko/stream/connectors/ftp/impl/SftpOperations.scala
@@ -191,7 +191,7 @@ private[ftp] trait SftpOperations { self: 
FtpLike[SSHClient, SftpSettings] =>
       }
     }
 
-  private[this] def authPublickey(identity: SftpIdentity)(implicit ssh: 
SSHClient) = {
+  private def authPublickey(identity: SftpIdentity)(implicit ssh: SSHClient) = 
{
     def bats(array: Array[Byte]): String = new String(array, 
StandardCharsets.UTF_8)
 
     val passphrase =
diff --git 
a/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeCQueryGraphLogic.scala
 
b/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeCQueryGraphLogic.scala
index 9dfb8764a..717a5dd27 100644
--- 
a/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeCQueryGraphLogic.scala
+++ 
b/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeCQueryGraphLogic.scala
@@ -44,7 +44,7 @@ private[geode] abstract class GeodeCQueryGraphLogic[V](val 
shape: SourceShape[V]
 
   val onElement: AsyncCallback[V]
 
-  private var query: CqQuery = _
+  private var query: CqQuery = null
 
   override def executeQuery() = Try {
 
diff --git 
a/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeSourceStageLogic.scala
 
b/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeSourceStageLogic.scala
index 00cc25364..5249198a7 100644
--- 
a/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeSourceStageLogic.scala
+++ 
b/geode/src/main/scala/org/apache/pekko/stream/connectors/geode/impl/stage/GeodeSourceStageLogic.scala
@@ -25,7 +25,7 @@ import scala.util.{ Failure, Success, Try }
 private[geode] abstract class GeodeSourceStageLogic[V](shape: SourceShape[V], 
clientCache: ClientCache)
     extends GraphStageLogic(shape) {
 
-  protected var initialResultsIterator: java.util.Iterator[V] = _
+  protected var initialResultsIterator: java.util.Iterator[V] = null
 
   val onConnect: AsyncCallback[Unit]
 
diff --git 
a/google-cloud-bigquery-storage/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/bigquery/storage/javadsl/BigQueryStorage.scala
 
b/google-cloud-bigquery-storage/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/bigquery/storage/javadsl/BigQueryStorage.scala
index 34bb82b2e..c2a6831a7 100644
--- 
a/google-cloud-bigquery-storage/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/bigquery/storage/javadsl/BigQueryStorage.scala
+++ 
b/google-cloud-bigquery-storage/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/bigquery/storage/javadsl/BigQueryStorage.scala
@@ -113,7 +113,7 @@ object BigQueryStorage {
       : Source[(ReadSession.Schema, 
java.util.List[Source[ReadRowsResponse.Rows, NotUsed]]), 
CompletionStage[NotUsed]] =
     create(projectId, datasetId, tableId, dataFormat, None, maxNumStreams)
 
-  private[this] def create(
+  private def create(
       projectId: String,
       datasetId: String,
       tableId: String,
@@ -206,7 +206,7 @@ object BigQueryStorage {
       um: Unmarshaller[ByteString, A]): Source[A, CompletionStage[NotUsed]] =
     createMergedStreams(projectId, datasetId, tableId, dataFormat, None, 
maxNumStreams, um.asScala)
 
-  private[this] def createMergedStreams[A](
+  private def createMergedStreams[A](
       projectId: String,
       datasetId: String,
       tableId: String,
diff --git 
a/google-cloud-bigquery/src/main/mima-filters/2.0.x.backward.excludes/context-bounds.excludes
 
b/google-cloud-bigquery/src/main/mima-filters/2.0.x.backward.excludes/context-bounds.excludes
new file mode 100644
index 000000000..a8b18fe12
--- /dev/null
+++ 
b/google-cloud-bigquery/src/main/mima-filters/2.0.x.backward.excludes/context-bounds.excludes
@@ -0,0 +1,22 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+# Converting context bound [T: FromResponseUnmarshaller] to explicit implicit
+# parameter reorders the implicit parameter list at the bytecode level.
+# This is a binary-incompatible change but does not affect source-level callers
+# using implicit resolution. BigQuery API is @ApiMayChange.
+ProblemFilters.exclude[IncompatibleMethTypeProblem]("org.apache.pekko.stream.connectors.googlecloud.bigquery.scaladsl.BigQuery.singleRequest")
diff --git 
a/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/impl/GCStorageStreamIntegrationSpec.scala
 
b/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/impl/GCStorageStreamIntegrationSpec.scala
index cf3337dd0..4f2e7b386 100644
--- 
a/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/impl/GCStorageStreamIntegrationSpec.scala
+++ 
b/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/impl/GCStorageStreamIntegrationSpec.scala
@@ -43,7 +43,7 @@ trait GCStorageStreamIntegrationSpec
   private implicit val defaultPatience: PatienceConfig =
     PatienceConfig(timeout = 60.seconds, interval = 60.millis)
 
-  var folderName: String = _
+  var folderName: String = null
 
   def testFileName(file: String): String = folderName + file
 
diff --git 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/http/GoogleHttp.scala
 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/http/GoogleHttp.scala
index d4b79af93..b2eb69b5f 100644
--- 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/http/GoogleHttp.scala
+++ 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/http/GoogleHttp.scala
@@ -87,12 +87,13 @@ private[connectors] final class GoogleHttp private (val 
http: HttpExt) extends A
    * If `authenticate = true` adds an Authorization header to each request.
    * Retries the request if the [[FromResponseUnmarshaller]] throws a 
[[pekko.stream.connectors.google.util.Retry]].
    */
-  def cachedHostConnectionPool[T: FromResponseUnmarshaller](
+  def cachedHostConnectionPool[T](
       host: String,
       port: Int = -1,
       https: Boolean = true,
       authenticate: Boolean = true,
-      parallelism: Int = 1): Flow[HttpRequest, T, Future[HostConnectionPool]] =
+      parallelism: Int = 1)(implicit um: FromResponseUnmarshaller[T])
+      : Flow[HttpRequest, T, Future[HostConnectionPool]] =
     Flow[HttpRequest]
       .map((_, ()))
       .viaMat(cachedHostConnectionPoolWithContext[T, Unit](host, port, https, 
authenticate, parallelism).asFlow)(
@@ -104,12 +105,13 @@ private[connectors] final class GoogleHttp private (val 
http: HttpExt) extends A
    * If `authenticate = true` adds an Authorization header to each request.
    * Retries the request if the [[FromResponseUnmarshaller]] throws a 
[[pekko.stream.connectors.google.util.Retry]].
    */
-  def cachedHostConnectionPoolWithContext[T: FromResponseUnmarshaller, Ctx](
+  def cachedHostConnectionPoolWithContext[T, Ctx](
       host: String,
       port: Int = -1,
       https: Boolean = true,
       authenticate: Boolean = true,
-      parallelism: Int = 1): FlowWithContext[HttpRequest, Ctx, Try[T], Ctx, 
Future[HostConnectionPool]] =
+      parallelism: Int = 1)(implicit um: FromResponseUnmarshaller[T])
+      : FlowWithContext[HttpRequest, Ctx, Try[T], Ctx, 
Future[HostConnectionPool]] =
     FlowWithContext.fromTuples {
       Flow.fromMaterializer { (mat, attr) =>
         implicit val settings: GoogleSettings = 
GoogleAttributes.resolveSettings(mat, attr)
diff --git 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/javadsl/Google.scala
 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/javadsl/Google.scala
index 223bfed42..836a10cc2 100644
--- 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/javadsl/Google.scala
+++ 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/javadsl/Google.scala
@@ -47,7 +47,7 @@ private[connectors] trait Google {
       unmarshaller: Unmarshaller[HttpResponse, T],
       settings: GoogleSettings,
       system: ClassicActorSystemProvider): CompletionStage[T] =
-    ScalaGoogle.singleRequest[T](request)(unmarshaller.asScala, system, 
settings).asJava
+    ScalaGoogle.singleRequest[T](request)(system, settings, 
unmarshaller.asScala).asJava
 
   /**
    * Makes a series of requests to page through a resource. Authentication is 
handled automatically.
diff --git 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/jwt/JwtSprayJson.scala
 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/jwt/JwtSprayJson.scala
index 5ae401f67..d07b213b3 100644
--- 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/jwt/JwtSprayJson.scala
+++ 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/jwt/JwtSprayJson.scala
@@ -82,9 +82,9 @@ private[google] class JwtSprayJson private (defaultClock: 
Clock)
       jwtId = safeGetField[String](jsObj, "jti"))
   }
 
-  private[this] def safeRead[A: JsonReader](js: JsValue) =
+  private def safeRead[A: JsonReader](js: JsValue) =
     safeReader[A].read(js).fold(_ => None, a => Option(a))
 
-  private[this] def safeGetField[A: JsonReader](js: JsObject, name: String) =
+  private def safeGetField[A: JsonReader](js: JsObject, name: String) =
     js.fields.get(name).flatMap(safeRead[A])
 }
diff --git 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/scaladsl/Google.scala
 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/scaladsl/Google.scala
index 0c7510e09..16c60d632 100644
--- 
a/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/scaladsl/Google.scala
+++ 
b/google-common/src/main/scala/org/apache/pekko/stream/connectors/google/scaladsl/Google.scala
@@ -40,8 +40,9 @@ private[connectors] trait Google {
    * @tparam T the data model for the resource
    * @return a [[scala.concurrent.Future]] containing the unmarshalled response
    */
-  final def singleRequest[T: FromResponseUnmarshaller](
-      request: HttpRequest)(implicit system: ClassicActorSystemProvider, 
settings: GoogleSettings): Future[T] =
+  final def singleRequest[T](
+      request: HttpRequest)(implicit system: ClassicActorSystemProvider, 
settings: GoogleSettings,
+      um: FromResponseUnmarshaller[T]): Future[T] =
     GoogleHttp().singleAuthenticatedRequest[T](request)
 
   /**
@@ -64,7 +65,8 @@ private[connectors] trait Google {
    * @tparam Out the data model for the resource
    * @return a [[pekko.stream.scaladsl.Sink]] that materializes a 
[[scala.concurrent.Future]] containing the unmarshalled resource
    */
-  final def resumableUpload[Out: FromResponseUnmarshaller](request: 
HttpRequest): Sink[ByteString, Future[Out]] =
+  final def resumableUpload[Out](request: HttpRequest)(
+      implicit um: FromResponseUnmarshaller[Out]): Sink[ByteString, 
Future[Out]] =
     ResumableUpload[Out](request)
 
 }
diff --git a/hdfs/src/test/scala/docs/scaladsl/HdfsReaderSpec.scala 
b/hdfs/src/test/scala/docs/scaladsl/HdfsReaderSpec.scala
index 101eefde7..a21da067e 100644
--- a/hdfs/src/test/scala/docs/scaladsl/HdfsReaderSpec.scala
+++ b/hdfs/src/test/scala/docs/scaladsl/HdfsReaderSpec.scala
@@ -39,7 +39,7 @@ class HdfsReaderSpec
     with BeforeAndAfterEach
     with LogCapturing {
 
-  private var hdfsCluster: MiniDFSCluster = _
+  private var hdfsCluster: MiniDFSCluster = null
   private val destination = "/tmp/pekko-connectors/"
 
   implicit val system: ActorSystem = ActorSystem()
diff --git a/hdfs/src/test/scala/docs/scaladsl/HdfsWriterSpec.scala 
b/hdfs/src/test/scala/docs/scaladsl/HdfsWriterSpec.scala
index 1bfbee788..79aafa51c 100644
--- a/hdfs/src/test/scala/docs/scaladsl/HdfsWriterSpec.scala
+++ b/hdfs/src/test/scala/docs/scaladsl/HdfsWriterSpec.scala
@@ -41,7 +41,7 @@ class HdfsWriterSpec
     with BeforeAndAfterEach
     with LogCapturing {
 
-  private var hdfsCluster: MiniDFSCluster = _
+  private var hdfsCluster: MiniDFSCluster = null
   private val destination = "/tmp/pekko-connectors/"
 
   implicit val system: ActorSystem = ActorSystem()
diff --git 
a/influxdb/src/main/scala/org/apache/pekko/stream/connectors/influxdb/impl/InfluxDbSourceStage.scala
 
b/influxdb/src/main/scala/org/apache/pekko/stream/connectors/influxdb/impl/InfluxDbSourceStage.scala
index 9738b29a1..5fb62ce35 100644
--- 
a/influxdb/src/main/scala/org/apache/pekko/stream/connectors/influxdb/impl/InfluxDbSourceStage.scala
+++ 
b/influxdb/src/main/scala/org/apache/pekko/stream/connectors/influxdb/impl/InfluxDbSourceStage.scala
@@ -56,7 +56,7 @@ private[influxdb] final class InfluxDbSourceLogic[T](clazz: 
Class[T],
     shape: SourceShape[T])
     extends InfluxDbBaseSourceLogic[T](influxDB, query, outlet, shape) {
 
-  var resultMapperHelper: PekkoConnectorsResultMapperHelper = _
+  var resultMapperHelper: PekkoConnectorsResultMapperHelper = null
 
   override def preStart(): Unit = {
     resultMapperHelper = new PekkoConnectorsResultMapperHelper
diff --git a/influxdb/src/test/scala/docs/scaladsl/FlowSpec.scala 
b/influxdb/src/test/scala/docs/scaladsl/FlowSpec.scala
index 246443dba..c43d02663 100644
--- a/influxdb/src/test/scala/docs/scaladsl/FlowSpec.scala
+++ b/influxdb/src/test/scala/docs/scaladsl/FlowSpec.scala
@@ -47,7 +47,7 @@ class FlowSpec
 
   final val DatabaseName = this.getClass.getSimpleName
 
-  implicit var influxDB: InfluxDB = _
+  implicit var influxDB: InfluxDB = null
 
   override protected def beforeAll(): Unit =
     influxDB = setupConnection(DatabaseName)
diff --git a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala 
b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
index 35c86547d..45ebea21b 100644
--- a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
+++ b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
@@ -41,7 +41,7 @@ class InfluxDbSourceSpec
 
   implicit val system: ActorSystem = ActorSystem()
 
-  implicit var influxDB: InfluxDB = _
+  implicit var influxDB: InfluxDB = null
 
   override protected def beforeAll(): Unit =
     influxDB = setupConnection(DatabaseName)
diff --git a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSpec.scala 
b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSpec.scala
index 23494263f..d437a6243 100644
--- a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSpec.scala
+++ b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSpec.scala
@@ -46,7 +46,7 @@ class InfluxDbSpec
 
   final val DatabaseName = this.getClass.getSimpleName
 
-  implicit var influxDB: InfluxDB = _
+  implicit var influxDB: InfluxDB = null
 
   // #define-class
   override protected def beforeAll(): Unit = {
diff --git 
a/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPullStage.scala
 
b/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPullStage.scala
index 2209d120d..34602616e 100644
--- 
a/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPullStage.scala
+++ 
b/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPullStage.scala
@@ -68,7 +68,7 @@ private[ironmq] final class IronMqPullStage(queueName: 
String, settings: IronMqS
       private var fetching: Boolean = false
 
       private var buffer: List[ReservedMessage] = List.empty
-      private var client: IronMqClient = _ // set in preStart
+      private var client: IronMqClient = null // set in preStart
 
       override def preStart(): Unit =
         client = IronMqClient(settings)(materializer.system, materializer)
diff --git 
a/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPushStage.scala
 
b/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPushStage.scala
index dab27b8e1..c801396c1 100644
--- 
a/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPushStage.scala
+++ 
b/ironmq/src/main/scala/org/apache/pekko/stream/connectors/ironmq/impl/IronMqPushStage.scala
@@ -50,7 +50,7 @@ private[ironmq] class IronMqPushStage(queueName: String, 
settings: IronMqSetting
 
       private var runningFutures: Int = 0
       private var exceptionFromUpstream: Option[Throwable] = None
-      private var client: IronMqClient = _ // set in preStart
+      private var client: IronMqClient = null // set in preStart
 
       override def preStart(): Unit = {
         super.preStart()
diff --git 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/Envelopes.scala
 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/Envelopes.scala
index 25c4a32b5..4c10cf60a 100644
--- 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/Envelopes.scala
+++ 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/Envelopes.scala
@@ -29,7 +29,7 @@ case class AckEnvelope private[jakartams] (message: 
jms.Message, private val jms
 
 case class TxEnvelope private[jakartams] (message: jms.Message, private val 
jmsSession: JmsSession) {
 
-  private[this] val commitPromise = Promise[() => Unit]()
+  private val commitPromise = Promise[() => Unit]()
 
   private[jakartams] val commitFuture: Future[() => Unit] = 
commitPromise.future
 
diff --git 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsBrowseStage.scala
 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsBrowseStage.scala
index 13d5db579..fc80c0756 100644
--- 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsBrowseStage.scala
+++ 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsBrowseStage.scala
@@ -36,10 +36,10 @@ private[jakartams] final class JmsBrowseStage(settings: 
JmsBrowseSettings, queue
     new GraphStageLogic(shape) with OutHandler {
       setHandler(out, this)
 
-      var connection: jms.Connection = _
-      var session: jms.Session = _
-      var browser: jms.QueueBrowser = _
-      var messages: java.util.Enumeration[jms.Message] = _
+      var connection: jms.Connection = null
+      var session: jms.Session = null
+      var browser: jms.QueueBrowser = null
+      var messages: java.util.Enumeration[jms.Message] = null
 
       override def preStart(): Unit = {
         val ackMode = settings.acknowledgeMode.mode
diff --git 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsConnector.scala
 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsConnector.scala
index a70e208f4..9ac16f637 100644
--- 
a/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsConnector.scala
+++ 
b/jakartams/src/main/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsConnector.scala
@@ -39,7 +39,7 @@ private[jakartams] trait JmsConnector[S <: JmsSession] {
 
   import JmsConnector._
 
-  implicit protected var ec: ExecutionContext = _
+  implicit protected var ec: ExecutionContext = null
 
   private var jmsSessions = Seq.empty[S]
 
@@ -53,7 +53,7 @@ private[jakartams] trait JmsConnector[S <: JmsSession] {
 
   private val connectionFailedCB = 
getAsyncCallback[Throwable](connectionFailed)
 
-  private var connectionStateQueue: 
SourceQueueWithComplete[InternalConnectionState] = _
+  private var connectionStateQueue: 
SourceQueueWithComplete[InternalConnectionState] = null
 
   private val connectionStateSourcePromise = 
Promise[Source[InternalConnectionState, NotUsed]]()
 
diff --git 
a/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/JmsSharedServerSpec.scala
 
b/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/JmsSharedServerSpec.scala
index 9ee886ede..9b96d65c2 100644
--- 
a/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/JmsSharedServerSpec.scala
+++ 
b/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/JmsSharedServerSpec.scala
@@ -22,8 +22,8 @@ import scala.util.Random
  * Creates a single server and connection factory which is shared for all 
tests.
  */
 abstract class JmsSharedServerSpec extends JmsSpec {
-  private var jmsBroker: EmbeddedActiveMQServer = _
-  private var connectionFactory: ConnectionFactory = _
+  private var jmsBroker: EmbeddedActiveMQServer = null
+  private var connectionFactory: ConnectionFactory = null
   private val SHARED_SERVER_ID: Int = 2
 
   override def beforeAll(): Unit = {
diff --git 
a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/Envelopes.scala 
b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/Envelopes.scala
index 6758e6b21..6711478b6 100644
--- a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/Envelopes.scala
+++ b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/Envelopes.scala
@@ -29,7 +29,7 @@ case class AckEnvelope private[jms] (message: jms.Message, 
private val jmsSessio
 
 case class TxEnvelope private[jms] (message: jms.Message, private val 
jmsSession: JmsSession) {
 
-  private[this] val commitPromise = Promise[() => Unit]()
+  private val commitPromise = Promise[() => Unit]()
 
   private[jms] val commitFuture: Future[() => Unit] = commitPromise.future
 
diff --git 
a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsBrowseStage.scala
 
b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsBrowseStage.scala
index 45a0078d6..c8f9d76f3 100644
--- 
a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsBrowseStage.scala
+++ 
b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsBrowseStage.scala
@@ -36,10 +36,10 @@ private[jms] final class JmsBrowseStage(settings: 
JmsBrowseSettings, queue: Dest
     new GraphStageLogic(shape) with OutHandler {
       setHandler(out, this)
 
-      var connection: jms.Connection = _
-      var session: jms.Session = _
-      var browser: jms.QueueBrowser = _
-      var messages: java.util.Enumeration[jms.Message] = _
+      var connection: jms.Connection = null
+      var session: jms.Session = null
+      var browser: jms.QueueBrowser = null
+      var messages: java.util.Enumeration[jms.Message] = null
 
       override def preStart(): Unit = {
         val ackMode = settings.acknowledgeMode.mode
diff --git 
a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsConnector.scala
 
b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsConnector.scala
index b2f3706b9..d6e16caee 100644
--- 
a/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsConnector.scala
+++ 
b/jms/src/main/scala/org/apache/pekko/stream/connectors/jms/impl/JmsConnector.scala
@@ -40,7 +40,7 @@ private[jms] trait JmsConnector[S <: JmsSession] {
 
   import JmsConnector._
 
-  implicit protected var ec: ExecutionContext = _
+  implicit protected var ec: ExecutionContext = null
 
   private var jmsSessions = Seq.empty[S]
 
@@ -54,7 +54,7 @@ private[jms] trait JmsConnector[S <: JmsSession] {
 
   private val connectionFailedCB = 
getAsyncCallback[Throwable](connectionFailed)
 
-  private var connectionStateQueue: 
SourceQueueWithComplete[InternalConnectionState] = _
+  private var connectionStateQueue: 
SourceQueueWithComplete[InternalConnectionState] = null
 
   private val connectionStateSourcePromise = 
Promise[Source[InternalConnectionState, NotUsed]]()
 
diff --git 
a/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/JmsSharedServerSpec.scala
 
b/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/JmsSharedServerSpec.scala
index 2d32fce5e..70c2edf04 100644
--- 
a/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/JmsSharedServerSpec.scala
+++ 
b/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/JmsSharedServerSpec.scala
@@ -22,8 +22,8 @@ import scala.util.Random
  * Creates a single server and connection factory which is shared for all 
tests.
  */
 abstract class JmsSharedServerSpec extends JmsSpec {
-  private var jmsBroker: JmsBroker = _
-  private var connectionFactory: ConnectionFactory = _
+  private var jmsBroker: JmsBroker = null
+  private var connectionFactory: ConnectionFactory = null
 
   override def beforeAll(): Unit = {
     jmsBroker = JmsBroker()
diff --git 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
index 3e0022905..69a822b36 100644
--- 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
+++ 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
@@ -72,9 +72,9 @@ private[kinesis] class KinesisSchedulerSourceStage(
 
     // We're transmitting backpressure from the Outlet to the Scheduler using 
a Semaphore instance
     // semaphore.acquire ~> callback ~> push downstream ~> semaphore.release
-    private[this] val backpressureSemaphore = new Semaphore(bufferSize)
-    private[this] val buffer = mutable.Queue.empty[CommittableRecord]
-    private[this] var schedulerOpt: Option[Scheduler] = None
+    private val backpressureSemaphore = new Semaphore(bufferSize)
+    private val buffer = mutable.Queue.empty[CommittableRecord]
+    private var schedulerOpt: Option[Scheduler] = None
 
     override def preStart(): Unit = {
       implicit val ec: ExecutionContext = executionContext(attributes)
diff --git 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSourceStage.scala
 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSourceStage.scala
index b8db4c5c2..d88cbdaf0 100644
--- 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSourceStage.scala
+++ 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSourceStage.scala
@@ -69,9 +69,9 @@ private[kinesis] class KinesisSourceStage(shardSettings: 
ShardSettings, amazonKi
 
       import shardSettings._
 
-      private[this] var currentShardIterator: String = _
-      private[this] val buffer = mutable.Queue.empty[Record]
-      private[this] var self: StageActor = _
+      private var currentShardIterator: String = null
+      private val buffer = mutable.Queue.empty[Record]
+      private var self: StageActor = null
 
       override def preStart(): Unit = {
         self = getStageActor(awaitingShardIterator)
@@ -144,19 +144,19 @@ private[kinesis] class KinesisSourceStage(shardSettings: 
ShardSettings, amazonKi
           log.warning("unexpected timer [{}]", other)
       }
 
-      private[this] val handleGetRecords: Try[GetRecordsResponse] => Unit = {
+      private val handleGetRecords: Try[GetRecordsResponse] => Unit = {
         case Failure(exception) => self.ref ! GetRecordsFailure(exception)
         case Success(result)    => self.ref ! GetRecordsSuccess(result)
       }
 
-      private[this] def requestRecords(): Unit =
+      private def requestRecords(): Unit =
         amazonKinesisAsync
           .getRecords(
             
GetRecordsRequest.builder().limit(limit).shardIterator(currentShardIterator).build())
           .asScala
           .onComplete(handleGetRecords)(parasitic)
 
-      private[this] def requestShardIterator(): Unit = {
+      private def requestShardIterator(): Unit = {
         val request = Function
           .chain[GetShardIteratorRequest.Builder](
             Seq(
diff --git 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/ShardProcessor.scala
 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/ShardProcessor.scala
index 704e17e14..3fd0f4fff 100644
--- 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/ShardProcessor.scala
+++ 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/ShardProcessor.scala
@@ -48,8 +48,8 @@ private[kinesis] class ShardProcessor(
   // by the AWS KCL Scheduler after the Scheduler run() method has been 
invoked.
   private val lastRecordSemaphore = new Semaphore(1)
 
-  private var shardData: ShardProcessorData = _
-  private var checkpointer: RecordProcessorCheckpointer = _
+  private var shardData: ShardProcessorData = null
+  private var checkpointer: RecordProcessorCheckpointer = null
   private var shutdown: Option[ShutdownReason] = None
 
   override def initialize(initializationInput: InitializationInput): Unit =
diff --git 
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisSchedulerSourceSpec.scala
 
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisSchedulerSourceSpec.scala
index 70b72b284..aafeb57d6 100644
--- 
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisSchedulerSourceSpec.scala
+++ 
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisSchedulerSourceSpec.scala
@@ -270,8 +270,8 @@ class KinesisSchedulerSourceSpec
 
     private val semaphore = new Semaphore(0)
 
-    var recordProcessor: ShardRecordProcessor = _
-    var otherRecordProcessor: ShardRecordProcessor = _
+    var recordProcessor: ShardRecordProcessor = null
+    var otherRecordProcessor: ShardRecordProcessor = null
     private val schedulerBuilder = { (x: ShardRecordProcessorFactory) =>
       recordProcessor = x.shardRecordProcessor()
       otherRecordProcessor = x.shardRecordProcessor()
@@ -334,7 +334,7 @@ class KinesisSchedulerSourceSpec
     "checkpoint batch of records with same sequence number" in new 
KinesisSchedulerCheckpointContext {
       val checkpointer: KinesisClientRecord => Unit =
         org.mockito.Mockito.mock(classOf[KinesisClientRecord => Unit])
-      var latestRecord: KinesisClientRecord = _
+      var latestRecord: KinesisClientRecord = null
       val allRecordsPushed: Future[Unit] = Future {
         for (i <- 1 to 3) {
           val clientRecord = 
org.mockito.Mockito.mock(classOf[KinesisClientRecord])
@@ -370,8 +370,8 @@ class KinesisSchedulerSourceSpec
     "checkpoint batch of records of different shards" in new 
KinesisSchedulerCheckpointContext {
       val checkpointerShard1: KinesisClientRecord => Unit =
         org.mockito.Mockito.mock(classOf[KinesisClientRecord => Unit])
-      var latestRecordShard1: KinesisClientRecord = _
-      var latestRecordShard2: KinesisClientRecord = _
+      var latestRecordShard1: KinesisClientRecord = null
+      var latestRecordShard2: KinesisClientRecord = null
       val checkpointerShard2: KinesisClientRecord => Unit =
         org.mockito.Mockito.mock(classOf[KinesisClientRecord => Unit])
 
diff --git 
a/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbFlowStage.scala
 
b/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbFlowStage.scala
index 6230f42f4..5f239aeee 100644
--- 
a/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbFlowStage.scala
+++ 
b/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbFlowStage.scala
@@ -55,8 +55,8 @@ private[orientdb] class OrientDbFlowStage[T, C](
 
   sealed abstract class OrientDbLogic extends GraphStageLogic(shape) with 
InHandler with OutHandler {
 
-    protected var client: ODatabaseSession = _
-    protected var oObjectClient: OObjectDatabaseTx = _
+    protected var client: ODatabaseSession = null
+    protected var oObjectClient: OObjectDatabaseTx = null
 
     override def preStart(): Unit = {
       client = settings.oDatabasePool.acquire()
diff --git 
a/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbSourceStage.scala
 
b/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbSourceStage.scala
index 0b91a2655..db7b85eec 100644
--- 
a/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbSourceStage.scala
+++ 
b/orientdb/src/main/scala/org/apache/pekko/stream/connectors/orientdb/impl/OrientDbSourceStage.scala
@@ -126,8 +126,8 @@ private[orientdb] final class 
OrientDbSourceStage[T](className: String,
 
   private abstract class Logic extends GraphStageLogic(shape) with OutHandler {
 
-    protected var client: ODatabaseSession = _
-    protected var oObjectClient: OObjectDatabaseTx = _
+    protected var client: ODatabaseSession = null
+    protected var oObjectClient: OObjectDatabaseTx = null
     protected var skip = settings.skip
 
     override def preStart(): Unit = {
diff --git a/orientdb/src/test/scala/docs/scaladsl/OrientDbSpec.scala 
b/orientdb/src/test/scala/docs/scaladsl/OrientDbSpec.scala
index c6953f620..d3122bc88 100644
--- a/orientdb/src/test/scala/docs/scaladsl/OrientDbSpec.scala
+++ b/orientdb/src/test/scala/docs/scaladsl/OrientDbSpec.scala
@@ -72,9 +72,9 @@ class OrientDbSpec extends AnyWordSpec with Matchers with 
BeforeAndAfterAll with
   case class Book(title: String)
   // #define-class
 
-  var orientDB: OrientDB = _
-  var oDatabase: ODatabasePool = _
-  var client: ODatabaseSession = _
+  var orientDB: OrientDB = null
+  var oDatabase: ODatabasePool = null
+  var client: ODatabaseSession = null
 
   override def beforeAll() = {
     orientDB = new OrientDB(url, username, password, 
OrientDBConfig.defaultConfig())
diff --git 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaFlow.scala
 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaFlow.scala
index d95e7872c..1cb4a0795 100644
--- 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaFlow.scala
+++ 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaFlow.scala
@@ -41,7 +41,7 @@ import scala.util.{ Failure, Success, Try }
 
   val clientConfig = writerSettings.clientConfig
 
-  private var writer: EventStreamWriter[A] = _
+  private var writer: EventStreamWriter[A] = null
 
   private val semaphore = new Semaphore(writerSettings.maximumInflightMessages)
 
diff --git 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaSource.scala
 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaSource.scala
index 8b0afc50d..9860ab763 100644
--- 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaSource.scala
+++ 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaSource.scala
@@ -47,7 +47,7 @@ import scala.util.{ Failure, Success, Try }
 
   private def out: Outlet[PravegaEvent[A]] = shape.out
 
-  private var reader: EventStreamReader[A] = _
+  private var reader: EventStreamReader[A] = null
 
   protected val clientConfig: ClientConfig = readerSettings.clientConfig
 
diff --git 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableReadFlow.scala
 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableReadFlow.scala
index 462d0f7c2..4c4bbe022 100644
--- 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableReadFlow.scala
+++ 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableReadFlow.scala
@@ -42,8 +42,8 @@ import scala.util.control.NonFatal
   private def in: Inlet[K] = shape.in
   private def out: Outlet[Option[V]] = shape.out
 
-  private var keyValueTableFactory: KeyValueTableFactory = _
-  private var table: KeyValueTable = _
+  private var keyValueTableFactory: KeyValueTableFactory = null
+  private var table: KeyValueTable = null
 
   @volatile
   private var inFlight = 0
diff --git 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableSource.scala
 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableSource.scala
index e5d6ed852..372fee83a 100644
--- 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableSource.scala
+++ 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableSource.scala
@@ -51,8 +51,8 @@ import io.pravega.common.util.AsyncIterator
 
   private def out = shape.out
 
-  private var keyValueTableFactory: KeyValueTableFactory = _
-  private var table: KeyValueTable = _
+  private var keyValueTableFactory: KeyValueTableFactory = null
+  private var table: KeyValueTable = null
 
   private val queue = mutable.Queue.empty[TableEntry[V]]
 
diff --git 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableWriteFlow.scala
 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableWriteFlow.scala
index 019746acc..35bb1420c 100644
--- 
a/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableWriteFlow.scala
+++ 
b/pravega/src/main/scala/org/apache/pekko/stream/connectors/pravega/impl/PravegaTableWriteFlow.scala
@@ -47,9 +47,9 @@ import scala.util.control.NonFatal
   private def in = shape.in
   private def out = shape.out
 
-  private var keyValueTableFactory: KeyValueTableFactory = _
+  private var keyValueTableFactory: KeyValueTableFactory = null
 
-  private var table: KeyValueTable = _
+  private var table: KeyValueTable = null
 
   private val onAir = new AtomicInteger
 
diff --git a/project/Common.scala b/project/Common.scala
index 61391a98a..4a9e38253 100644
--- a/project/Common.scala
+++ b/project/Common.scala
@@ -62,14 +62,27 @@ object Common extends AutoPlugin {
       "-feature",
       "-unchecked",
       "-deprecation",
-      "-Xlint",
-      "-Ywarn-dead-code",
-      "-Wconf:cat=unused-nowarn:s",
-      "-Wconf:msg=Prefer the Scala annotation over Java's `@Deprecated`:s",
       "-release:17"),
     scalacOptions ++= {
-      if (isScala3.value) Seq("-Yfuture-lazy-vals")
-      else Seq.empty
+      if (isScala3.value) Seq(
+        "-Yfuture-lazy-vals",
+        "-Wconf:msg=Implicit parameters should be provided with a `using` 
clause:s",
+        "-Wconf:msg=is deprecated for wildcard arguments of types:s",
+        "-Wconf:msg=The trailing ` _` for eta-expansion is unnecessary:s",
+        "-Wconf:msg=with as a type operator has been deprecated:s",
+        "-Wconf:msg=Unreachable case except for null:s",
+        "-Wconf:msg=is no longer supported for vararg splices:s",
+        "-Wconf:msg=is not declared infix:s",
+        "-Wconf:msg=auto insertion will be deprecated:s",
+        "-Wconf:msg=Ignoring \\[this\\] qualifier:s",
+        "-Wconf:msg=trait App in package scala is deprecated:s",
+        "-Wconf:msg=pattern binding uses refutable extractor:s",
+        "-Wconf:msg=bad option.*-Yfuture-lazy-vals:s")
+      else Seq(
+        "-Xlint",
+        "-Ywarn-dead-code",
+        "-Wconf:cat=unused-nowarn:s",
+        "-Wconf:msg=Prefer the Scala annotation over Java's `@Deprecated`:s")
     },
     Compile / doc / scalacOptions := scalacOptions.value ++ Seq(
       "-doc-title",
diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/HttpRequests.scala
 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/HttpRequests.scala
index 5e60d202f..d5bc46340 100644
--- 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/HttpRequests.scala
+++ 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/HttpRequests.scala
@@ -327,7 +327,7 @@ import scala.xml.NodeSeq
       .withDefaultHeaders(allHeaders)
   }
 
-  private[this] def s3Request(s3Location: S3Location, method: HttpMethod, 
uriFn: Uri => Uri = identity)(
+  private def s3Request(s3Location: S3Location, method: HttpMethod, uriFn: Uri 
=> Uri = identity)(
       implicit conf: S3Settings): HttpRequest = {
     val loc = s3Location.validate(conf)
     val s3RequestUri = uriFn(requestUri(loc.bucket, Some(loc.key)))
@@ -340,7 +340,7 @@ import scala.xml.NodeSeq
   }
 
   @throws(classOf[IllegalUriException])
-  private[this] def requestAuthority(bucket: String, region: Region)(implicit 
conf: S3Settings): Authority = {
+  private def requestAuthority(bucket: String, region: Region)(implicit conf: 
S3Settings): Authority = {
     (conf.endpointUrl, conf.accessStyle) match {
       case (None, PathAccessStyle) =>
         Authority(Uri.Host(s"s3.$region.amazonaws.com"))
@@ -356,7 +356,7 @@ import scala.xml.NodeSeq
     }
   }
 
-  private[this] def requestUri(bucket: String, key: Option[String])(implicit 
conf: S3Settings): Uri = {
+  private def requestUri(bucket: String, key: Option[String])(implicit conf: 
S3Settings): Uri = {
     val basePath = conf.accessStyle match {
       case PathAccessStyle =>
         Uri.Path / bucket
diff --git 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/Signer.scala
 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/Signer.scala
index fd7404659..ac26cb15a 100644
--- 
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/Signer.scala
+++ 
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/auth/Signer.scala
@@ -46,16 +46,16 @@ import pekko.stream.scaladsl.Source
         .mapMaterializedValue(_ => NotUsed)
     }
 
-  private[this] def sessionHeader(key: SigningKey): Option[HttpHeader] =
+  private def sessionHeader(key: SigningKey): Option[HttpHeader] =
     key.sessionToken.map(RawHeader("X-Amz-Security-Token", _))
 
-  private[this] def authorizationHeader(algorithm: String,
+  private def authorizationHeader(algorithm: String,
       key: SigningKey,
       requestDate: ZonedDateTime,
       canonicalRequest: CanonicalRequest): HttpHeader =
     RawHeader("Authorization", authorizationString(algorithm, key, 
requestDate, canonicalRequest))
 
-  private[this] def authorizationString(algorithm: String,
+  private def authorizationString(algorithm: String,
       key: SigningKey,
       requestDate: ZonedDateTime,
       canonicalRequest: CanonicalRequest): String = {
diff --git 
a/sns/src/test/scala/org/apache/pekko/stream/connectors/sns/IntegrationTestContext.scala
 
b/sns/src/test/scala/org/apache/pekko/stream/connectors/sns/IntegrationTestContext.scala
index 0cf134696..46f2be80a 100644
--- 
a/sns/src/test/scala/org/apache/pekko/stream/connectors/sns/IntegrationTestContext.scala
+++ 
b/sns/src/test/scala/org/apache/pekko/stream/connectors/sns/IntegrationTestContext.scala
@@ -35,8 +35,8 @@ trait IntegrationTestContext extends BeforeAndAfterAll with 
ScalaFutures {
 
   def snsEndpoint: String = s"http://localhost:4100";
 
-  implicit var snsClient: SnsAsyncClient = _
-  var topicArn: String = _
+  implicit var snsClient: SnsAsyncClient = null
+  var topicArn: String = null
 
   private val topicNumber = new AtomicInteger()
 
diff --git a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala 
b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
index e1a24c2f9..4fb4baee0 100644
--- a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
+++ b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
@@ -47,9 +47,9 @@ class SolrSpec extends AnyWordSpec with Matchers with 
BeforeAndAfterAll with Sca
 
   override implicit val patienceConfig: PatienceConfig = 
PatienceConfig(5.seconds)
 
-  private var cluster: MiniSolrCloudCluster = _
+  private var cluster: MiniSolrCloudCluster = null
 
-  private var zkTestServer: ZkTestServer = _
+  private var zkTestServer: ZkTestServer = null
   implicit val system: ActorSystem = ActorSystem()
   implicit val commitExecutionContext: ExecutionContext = 
ExecutionContext.global
 
diff --git 
a/sqs/src/main/scala/org/apache/pekko/stream/connectors/sqs/impl/BalancingMapAsync.scala
 
b/sqs/src/main/scala/org/apache/pekko/stream/connectors/sqs/impl/BalancingMapAsync.scala
index 8ec9b2893..157c00681 100644
--- 
a/sqs/src/main/scala/org/apache/pekko/stream/connectors/sqs/impl/BalancingMapAsync.scala
+++ 
b/sqs/src/main/scala/org/apache/pekko/stream/connectors/sqs/impl/BalancingMapAsync.scala
@@ -60,7 +60,7 @@ import scala.util.{ Failure, Success }
     new GraphStageLogic(shape) with InHandler with OutHandler {
 
       lazy val decider = 
inheritedAttributes.mandatoryAttribute[SupervisionStrategy].decider
-      var buffer: Buffer[Holder[Out]] = _
+      var buffer: Buffer[Holder[Out]] = null
       var parallelism = maxParallelism
 
       private val futureCB = getAsyncCallback[Holder[Out]](holder =>
diff --git 
a/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpBind.scala 
b/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpBind.scala
index 2eea60802..5bfc72b03 100644
--- 
a/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpBind.scala
+++ 
b/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpBind.scala
@@ -40,7 +40,7 @@ import scala.concurrent.{ Future, Promise }
   private def in = shape.in
   private def out = shape.out
 
-  private var listener: ActorRef = _
+  private var listener: ActorRef = null
 
   override def preStart(): Unit = {
     implicit val sender: ActorRef = getStageActor(processIncoming).ref
diff --git 
a/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpSend.scala 
b/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpSend.scala
index 038be1d4f..d12a2e7d3 100644
--- 
a/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpSend.scala
+++ 
b/udp/src/main/scala/org/apache/pekko/stream/connectors/udp/impl/UdpSend.scala
@@ -38,7 +38,7 @@ import scala.collection.immutable.Iterable
   private def in = shape.in
   private def out = shape.out
 
-  private var simpleSender: ActorRef = _
+  private var simpleSender: ActorRef = null
 
   override def preStart(): Unit = {
     getStageActor(processIncoming)
diff --git 
a/xml/src/main/scala/org/apache/pekko/stream/connectors/xml/impl/StreamingXmlParser.scala
 
b/xml/src/main/scala/org/apache/pekko/stream/connectors/xml/impl/StreamingXmlParser.scala
index 6d8496820..2ae32af6f 100644
--- 
a/xml/src/main/scala/org/apache/pekko/stream/connectors/xml/impl/StreamingXmlParser.scala
+++ 
b/xml/src/main/scala/org/apache/pekko/stream/connectors/xml/impl/StreamingXmlParser.scala
@@ -66,7 +66,7 @@ private[xml] object StreamingXmlParser {
   override def createLogic(inheritedAttributes: Attributes): GraphStageLogic =
     new GraphStageLogic(shape) with InHandler with OutHandler {
       private var started: Boolean = false
-      private var context: Ctx = _
+      private var context: Ctx = null.asInstanceOf[Ctx]
 
       import javax.xml.stream.XMLStreamConstants
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to