Skip to content
Open
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 @@ -25,7 +25,6 @@ import org.apache.pekko.stream.{ Materializer, SystemMaterializer }
import com.amazonaws.services.dynamodbv2.model._
import com.typesafe.config.Config

import scala.collection.immutable
import scala.concurrent.{ ExecutionContext, Future, Promise }
import scala.util.{ Success, Try }

Expand Down Expand Up @@ -107,7 +106,7 @@ class DynamoDBJournal(config: Config)
private case class OpFinished(pid: String, f: Future[Done])
private val opQueue: JMap[String, Future[Done]] = new JHMap

override def asyncWriteMessages(messages: immutable.Seq[AtomicWrite]): Future[immutable.Seq[Try[Unit]]] = {
override def asyncWriteMessages(messages: Seq[AtomicWrite]): Future[Seq[Try[Unit]]] = {
val p = Promise[Done]()
val pid = messages.head.persistenceId
opQueue.put(pid, p.future)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ object DynamoDBRecovery {
* @param partitionSeqNum - the partition sequence number for the given persistence id.
* @param partitionEventNums - will be 0-99, representing the event ordering within the given partition sequence.
*/
case class PartitionKeys(partitionSeqNum: Long, partitionEventNums: immutable.Seq[Long])
case class PartitionKeys(partitionSeqNum: Long, partitionEventNums: Seq[Long])

/**
* Groups Longs from a stream into a [PartitionKeys] whereas each sequence shall contain the values that would be within the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ class DynamoPartitionGroupedSpec extends TestKit(ActorSystem("DynamoPartitionGro
.request(1)
.expectNext(PartitionKeys(2L, 200L to 299))
.request(1)
.expectNext(PartitionKeys(3L, scala.collection.immutable.Seq(300L)))
.expectNext(PartitionKeys(3L, Seq(300L)))
.expectComplete()
}

Expand Down