Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference

import scala.annotation.tailrec
import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
Expand Down Expand Up @@ -584,10 +583,10 @@ class CassandraProjectionSpec
val entityId = UUID.randomUUID().toString
val projectionId = genRandomProjectionId()

def groupedHandler(): Handler[immutable.Seq[Envelope]] = new Handler[immutable.Seq[Envelope]] {
def groupedHandler(): Handler[Seq[Envelope]] = new Handler[Seq[Envelope]] {
private var state: Future[Option[ConcatStr]] = repository.findById(entityId)

override def process(group: immutable.Seq[Envelope]): Future[Done] = {
override def process(group: Seq[Envelope]): Future[Done] = {
val newState = state.flatMap { s =>
val concatStr = group.foldLeft(s) {
case (None, env) => Some(ConcatStr(env.id, env.message))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.cassandra.scaladsl

import scala.collection.immutable
import scala.concurrent.Future
import scala.concurrent.duration.Duration

Expand Down Expand Up @@ -88,7 +87,7 @@ object CassandraProjection {
def groupedWithin[Offset, Envelope](
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
handler: () => Handler[immutable.Seq[Envelope]]): GroupedProjection[Offset, Envelope] =
handler: () => Handler[Seq[Envelope]]): GroupedProjection[Offset, Envelope] =
new CassandraProjectionImpl(
projectionId,
sourceProvider,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,6 @@ import java.nio.charset.StandardCharsets
import java.util.Base64
import java.util.UUID

import scala.collection.immutable

import org.apache.pekko
import pekko.actor.ExtendedActorSystem
import pekko.actor.testkit.typed.scaladsl.LogCapturing
Expand Down Expand Up @@ -140,7 +138,7 @@ class OffsetSerializationSpec
}

val storageRepresentation = MultipleOffsets(
immutable.Seq(SingleOffset(ProjectionId(id.name, surrogateProjectionKey), LongManifest, "1", mergeable = true)))
Seq(SingleOffset(ProjectionId(id.name, surrogateProjectionKey), LongManifest, "1", mergeable = true)))

actualRep shouldBe storageRepresentation

Expand All @@ -155,7 +153,7 @@ class OffsetSerializationSpec
val mergeableOffset =
MergeableOffset(Map(surrogateProjectionKey1 -> 1L, surrogateProjectionKey2 -> 2L))
val storageRepresentation = MultipleOffsets(
immutable.Seq(
Seq(
SingleOffset(ProjectionId(projectionName, surrogateProjectionKey1), LongManifest, "1", mergeable = true),
SingleOffset(ProjectionId(projectionName, surrogateProjectionKey2), LongManifest, "2", mergeable = true)))

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ package org.apache.pekko.projection.internal.metrics.tools

import java.util.UUID

import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
Expand Down Expand Up @@ -161,8 +160,8 @@ object InternalProjectionStateMetricsSpec {
}
case groupedHandlerStrategy: GroupedHandlerStrategy[Envelope] @unchecked => {
val adaptedHandler = () =>
new Handler[immutable.Seq[Envelope]] {
override def process(envelopes: immutable.Seq[Envelope]): Future[Done] =
new Handler[Seq[Envelope]] {
override def process(envelopes: Seq[Envelope]): Future[Done] =
groupedHandlerStrategy.handlerFactory().process(envelopes).flatMap { _ =>
offsetStore.saveOffset(projectionId, envelopes.last.offset)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.internal.metrics.tools

import scala.collection.immutable
import scala.concurrent.Future

import org.apache.pekko
Expand Down Expand Up @@ -71,13 +70,13 @@ object TestHandlers {
* trigger an error and then be removed from the stack. To fail an item multiple times
* add its offset repeatedly. Uses `Int` instead of `Long` for convenience.
*/
def groupedWithErrors(erroredOffsets: Int*): () => Handler[immutable.Seq[Envelope]] = {
def groupedWithErrors(erroredOffsets: Int*): () => Handler[Seq[Envelope]] = {
var nextProcessStrategy = ProcessStrategy(erroredOffsets.map {
_.toLong
}.toList)
() =>
new Handler[immutable.Seq[Envelope]] {
override def process(envelopes: immutable.Seq[Envelope]): Future[Done] = {
new Handler[Seq[Envelope]] {
override def process(envelopes: Seq[Envelope]): Future[Done] = {
nextProcessStrategy match {
case SomeFailures(nextFail :: tail)
if envelopes
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.internal

import scala.collection.immutable
import scala.concurrent.Future
import scala.jdk.CollectionConverters._
import scala.jdk.FutureConverters._
Expand Down Expand Up @@ -55,13 +54,13 @@ import pekko.projection.scaladsl
}

/**
* INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to `scaladsl.Handler[immutable.Seq[Envelope]]`
* INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to `scaladsl.Handler[Seq[Envelope]]`
*/
@InternalApi private[projection] class GroupedHandlerAdapter[Envelope](
delegate: javadsl.Handler[java.util.List[Envelope]])
extends scaladsl.Handler[immutable.Seq[Envelope]] {
extends scaladsl.Handler[Seq[Envelope]] {

override def process(envelopes: immutable.Seq[Envelope]): Future[Done] = {
override def process(envelopes: Seq[Envelope]): Future[Done] = {
delegate.process(envelopes.asJava).asScala
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.internal

import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.Promise
Expand Down Expand Up @@ -89,7 +88,7 @@ private[projection] abstract class InternalProjectionState[Offset, Envelope](

protected def saveOffsetsAndReport(
projectionId: ProjectionId,
batch: immutable.Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
batch: Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {

// The batch contains multiple projections contexts. Each of these contexts may represent
// a single envelope or a group of envelopes. The size of the batch and the size of the
Expand All @@ -103,15 +102,15 @@ private[projection] abstract class InternalProjectionState[Offset, Envelope](
/**
* A convenience method to serialize asynchronous operations to occur one after another is complete
*/
private def serialize(batches: Map[String, immutable.Seq[ProjectionContextImpl[Offset, Envelope]]])(
op: (String, immutable.Seq[ProjectionContextImpl[Offset, Envelope]]) => Future[Done]): Future[Done] = {
private def serialize(batches: Map[String, Seq[ProjectionContextImpl[Offset, Envelope]]])(
op: (String, Seq[ProjectionContextImpl[Offset, Envelope]]) => Future[Done]): Future[Done] = {

val logProgressEvery: Int = 5
val size = batches.size
logger.debug("Processing [{}] partitioned batches serially", size)

def loop(
remaining: List[(String, immutable.Seq[ProjectionContextImpl[Offset, Envelope]])],
remaining: List[(String, Seq[ProjectionContextImpl[Offset, Envelope]])],
n: Int): Future[Done] = {
remaining match {
case Nil => Future.successful(Done)
Expand Down Expand Up @@ -247,11 +246,11 @@ private[projection] abstract class InternalProjectionState[Offset, Envelope](
HandlerRecoveryImpl[Offset, Envelope](projectionId, recoveryStrategy, logger, statusObserver, telemetry)

def processGrouped(
handler: Handler[immutable.Seq[Envelope]],
handler: Handler[Seq[Envelope]],
handlerRecovery: HandlerRecoveryImpl[Offset, Envelope],
envelopesAndOffsets: immutable.Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
envelopesAndOffsets: Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {

def processEnvelopes(partitioned: immutable.Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
def processEnvelopes(partitioned: Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
val first = partitioned.head
val firstOffset = first.offset
val lastOffset = partitioned.last.offset
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,6 @@ package org.apache.pekko.projection.internal
import java.util.Base64
import java.util.UUID

import scala.collection.immutable

import org.apache.pekko
import pekko.actor.typed.ActorSystem
import pekko.annotation.InternalApi
Expand All @@ -34,7 +32,7 @@ import pekko.serialization.Serializers
sealed trait StorageRepresentation
final case class SingleOffset(id: ProjectionId, manifest: String, offsetStr: String, mergeable: Boolean = false)
extends StorageRepresentation
final case class MultipleOffsets(reps: immutable.Seq[SingleOffset]) extends StorageRepresentation
final case class MultipleOffsets(reps: Seq[SingleOffset]) extends StorageRepresentation

final val StringManifest = "STR"
final val LongManifest = "LNG"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.internal

import scala.collection.immutable
import scala.concurrent.duration.FiniteDuration

import org.apache.pekko
Expand Down Expand Up @@ -128,10 +127,10 @@ private[projection] final case class SingleHandlerStrategy[Envelope](handlerFact
*/
@InternalApi
private[projection] final case class GroupedHandlerStrategy[Envelope](
handlerFactory: () => Handler[immutable.Seq[Envelope]],
handlerFactory: () => Handler[Seq[Envelope]],
afterEnvelopes: Option[Int] = None,
orAfterDuration: Option[FiniteDuration] = None)
extends FunctionHandlerStrategy[immutable.Seq[Envelope]](handlerFactory)
extends FunctionHandlerStrategy[Seq[Envelope]](handlerFactory)

/**
* INTERNAL API
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ package org.apache.pekko.projection.internal

import java.util

import scala.collection.immutable
import scala.jdk.CollectionConverters._

import org.apache.pekko
Expand Down Expand Up @@ -122,7 +121,7 @@ trait Telemetry {
dynamicAccess
.createInstanceFor[Telemetry](
fqcn,
immutable.Seq((classOf[ProjectionId], projectionId), (classOf[ActorSystem[?]], system)))
Seq((classOf[ProjectionId], projectionId), (classOf[ActorSystem[?]], system)))
.get
}
}
Expand Down
2 changes: 1 addition & 1 deletion docs/src/main/paradox/cassandra.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ The envelopes are grouped within a time window, or limited by a number of envelo
This window can be defined with `withGroup` of the returned `GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.

When using `groupedWithin` the handler is a @scala[`Handler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`Handler<List<EventEnvelope<ShoppingCart.Event>>>`].
When using `groupedWithin` the handler is a @scala[`Handler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`Handler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).

It stores the offset in Cassandra immediately after the `handler` has processed the envelopes, but that
Expand Down
2 changes: 1 addition & 1 deletion docs/src/main/paradox/jdbc.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ The envelopes are grouped within a time window, or limited by a number of envelo
This window can be defined with `withGroup` of the returned `GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.

When using `groupedWithin` the handler is a @scala[`JdbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`JdbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
When using `groupedWithin` the handler is a @scala[`JdbcHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`JdbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).

The offset is stored in the same transaction used for the user defined `handler`, which means exactly-once
Expand Down
2 changes: 1 addition & 1 deletion docs/src/main/paradox/r2dbc.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ The envelopes are grouped within a time window, or limited by a number of envelo
This window can be defined with `withGroup` of the returned `GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.

When using `groupedWithin` the handler is a @scala[`R2dbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`R2dbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
When using `groupedWithin` the handler is a @scala[`R2dbcHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`R2dbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).

The offset is stored in the same transaction used for the user defined `handler`, which means exactly-once
Expand Down
2 changes: 1 addition & 1 deletion docs/src/main/paradox/slick.md
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ The envelopes are grouped within a time window, or limited by a number of envelo
This window can be defined with `withGroup` of the returned `GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.

When using `groupedWithin` the handler is a `SlickHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`.
When using `groupedWithin` the handler is a `SlickHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`.
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).

The offset is stored in the same transaction as the `DBIO` returned from the `handler`, which means exactly-once
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.projection.state.scaladsl

import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future

Expand Down Expand Up @@ -112,7 +111,7 @@ object DurableStateSourceProvider {
def sliceRanges(
system: ActorSystem[?],
durableStateStoreQueryPluginId: String,
numberOfRanges: Int): immutable.Seq[Range] =
numberOfRanges: Int): Seq[Range] =
DurableStateStoreRegistry(system)
.durableStateStoreFor[DurableStateStoreBySliceQuery[Any]](durableStateStoreQueryPluginId)
.sliceRanges(numberOfRanges)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ package org.apache.pekko.projection.eventsourced.scaladsl

import java.time.Instant

import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import com.typesafe.config.Config
Expand Down Expand Up @@ -141,15 +140,15 @@ object EventSourcedProvider {
.readJournalFor[EventsBySliceQuery](readJournalPluginId, readJournalConfig)
.sliceForPersistenceId(persistenceId)

def sliceRanges(system: ActorSystem[?], readJournalPluginId: String, numberOfRanges: Int): immutable.Seq[Range] =
def sliceRanges(system: ActorSystem[?], readJournalPluginId: String, numberOfRanges: Int): Seq[Range] =
PersistenceQuery(system).readJournalFor[EventsBySliceQuery](readJournalPluginId).sliceRanges(numberOfRanges)

/** @since 2.0.0 */
def sliceRanges(
system: ActorSystem[?],
readJournalPluginId: String,
readJournalConfig: Config,
numberOfRanges: Int): immutable.Seq[Range] =
numberOfRanges: Int): Seq[Range] =
PersistenceQuery(system)
.readJournalFor[EventsBySliceQuery](readJournalPluginId, readJournalConfig)
.sliceRanges(numberOfRanges)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@

package org.apache.pekko.projection.eventsourced.scaldsl

import scala.collection.immutable.Seq
import scala.concurrent.Future
import com.typesafe.config.ConfigFactory
import org.apache.pekko
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -99,15 +99,13 @@ object JdbcProjectionDocExample {
// #handler

// #grouped-handler
import scala.collection.immutable

class GroupedShoppingCartHandler(repository: OrderRepository)
extends JdbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]], PlainJdbcSession] {
extends JdbcHandler[Seq[EventEnvelope[ShoppingCart.Event]], PlainJdbcSession] {
private val logger = LoggerFactory.getLogger(getClass)

override def process(
session: PlainJdbcSession,
envelopes: immutable.Seq[EventEnvelope[ShoppingCart.Event]]): Unit = {
envelopes: Seq[EventEnvelope[ShoppingCart.Event]]): Unit = {

// save all events in DB
envelopes.map(_.event).foreach {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,13 +94,11 @@ class SlickProjectionDocExample {
// #handler

// #grouped-handler
import scala.collection.immutable

class GroupedShoppingCartHandler(repository: OrderRepository)(implicit ec: ExecutionContext)
extends SlickHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]] {
extends SlickHandler[Seq[EventEnvelope[ShoppingCart.Event]]] {
private val logger = LoggerFactory.getLogger(getClass)

override def process(envelopes: immutable.Seq[EventEnvelope[ShoppingCart.Event]]): DBIO[Done] = {
override def process(envelopes: Seq[EventEnvelope[ShoppingCart.Event]]): DBIO[Done] = {
val dbios = envelopes.map(_.event).map {
case ShoppingCart.CheckedOut(cartId, time) =>
logger.info(s"Shopping cart $cartId was checked out at $time")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,6 @@ import org.scalatest.wordspec.AnyWordSpecLike

import java.time.Instant
import java.util.concurrent.ConcurrentHashMap
import scala.collection.immutable
import scala.concurrent.Future
import scala.concurrent.Promise

Expand Down Expand Up @@ -92,7 +91,7 @@ object EventProducerServiceSpec {
override def sliceForPersistenceId(persistenceId: String): Int =
persistenceExt.sliceForPersistenceId(persistenceId)

override def sliceRanges(numberOfRanges: Int): immutable.Seq[Range] =
override def sliceRanges(numberOfRanges: Int): Seq[Range] =
persistenceExt.sliceRanges(numberOfRanges)
}
}
Expand Down
Loading