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]