This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-persistence-r2dbc.git
The following commit(s) were added to refs/heads/main by this push:
new 98fd2c9 feat: add opt-in batched write journal to coalesce concurrent
writes … (#489)
98fd2c9 is described below
commit 98fd2c961613d757004e8ad45deed72debdbff82
Author: Jabir S. Minjibir <[email protected]>
AuthorDate: Wed Sep 30 05:55:07 2026 -0400
feat: add opt-in batched write journal to coalesce concurrent writes …
(#489)
* feat: add opt-in batched write journal to coalesce concurrent writes
(#471)
Motivation:
Under high concurrency across many persistence IDs, the default journal
issues
separate statements and commits for each write. This increases database
round-trips
and connection pool contention.
Modification:
- Add R2dbcBatchJournal to queue writes and flush them in batches.
- Add max-batch-size and max-batch-time configuration settings.
- Add R2dbcBatchJournalSpec to verify TCK compliance.
- Add R2dbcBatchJournalPerfSpec and R2dbcBatchJournalPerfManyActorsSpec for
performance testing.
Result:
Users can opt in to batched writes across persistence IDs for higher write
throughput.
Tests:
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournalSpec"
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournalPerfSpec"
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournalPerfManyActorsSpec"
- sbt +mimaReportBinaryIssues
- sbt validatePullRequest
References:
Fixes #471
* fix: address review feedback on batched write journal (#471)
Motivation:
Review feedback on #489 flagged blockers and cleanups. This commit also
restores an incremental review history by keeping the original commit
intact and applying the fixes on top.
Modification:
- Bound the write queue with max-queue-size and reject writes on overflow
(A1).
- Complete batch promises before publishing (A2); remove the implicit
Array-to-Seq copy from the flush hot path (A3).
- Require the postgres or yugabyte dialect at startup (A4).
- Move the plugin to its own batched-journal config id that inherits from
journal, so it can run side by side with the default journal and the
batch keys no longer leak into the default journal settings (A5, B5).
- Handle Flush without dead letters (B1); suppress WriteFinished and
FlushDone sends during shutdown (C).
- Restore the Lightbend copyright header on the derived file and revert all
changes to existing classes and local build files so the PR is purely
additive.
- Add fail-fast and queue-limit specs, plus docs for the plugin id, the
worst-case bisection cost, and the per-request event count.
Result:
The PR diff against main is insertions only: new class, new specs, one
additive reference.conf block, and docs. Reviewers can diff 87df38e..HEAD
to see just this delta.
Tests:
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" (75 tests, live
PostgreSQL)
- sbt +headerCheckAll
- sbt scalafmtCheckAll
- sbt +mimaReportBinaryIssues
- sbt docs/paradox
- validatePullRequest task not present on this branch; equivalents run
manually
References:
Refs #471
* fix: restore remaining upstream files so PR is purely additive (#471)
Motivation:
The original commit edited three files outside the batched journal: it
removed the yugabyte/yugabyte-db#10995 FIXME link from R2dbcSettings,
rewrote comments in JournalDao, and added a settings test to
R2dbcSettingsSpec. The review-fix commit message promised all changes to
existing classes were reverted; this commit completes that promise.
Modification:
Restore R2dbcSettings.scala, JournalDao.scala, and R2dbcSettingsSpec.scala
to their upstream content. The mixed-persistence-id contract already lives
in the R2dbcBatchJournal class scaladoc, so no documentation is lost.
Result:
The pull request diff against main is insertions only: new class, new
specs, one additive reference.conf block, docs, and a MultiPluginSpec
addition demonstrating side-by-side plugin use.
Tests:
- sbt "core / Test / compile"
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.R2dbcSettingsSpec"
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*"
- git diff --check
References:
Refs #471
* fix: use standard header for R2dbcBatchJournal (#471)
Motivation:
The class carried the Lightbend copyright block copied from
R2dbcJournal. New classes must use the standard Apache header.
Modification:
Remove the Lightbend copyright block from R2dbcBatchJournal. The
tests in this PR already use the standard header.
Result:
All new files in this PR use the standard header.
Tests:
- sbt headerCheckAll
References:
Refs #471
* fix: skip batched journal specs on unsupported dialects (#471)
Motivation:
The batched journal requires the postgres or yugabyte dialect, but
the MySQL CI job runs all core tests. The suites would fail at
journal startup with the dialect require. Two specs create the
journal in the class body, failing during suite construction.
Modification:
- Add a BatchedJournalDialectGate trait that pends tests when the
configured dialect is mysql, and mix it into all batched journal
suites except R2dbcBatchJournalMysqlDialectSpec, which expects the
dialect require.
- Make the test journal fields lazy so the dialect require can not
fail during suite construction.
Result:
The batched journal suites run on postgres and yugabyte and pend on
the mysql CI job.
Tests:
- sbt "core/Test/testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" (75 tests, live
PostgreSQL)
- sbt "++3.3.8 core/Test/compile"
- sbt +headerCheckAll
- scalafmt --mode diff-ref=origin/main
References:
Refs #471
* fix: address second review round on batched write journal (#471)
Standard ASF headers on all new files, import cleanup, Option.when,
math.min and Future.failed simplifications, atomicWrite returns
Try[Seq[SerializedJournalRow]], and the experimental note in the docs.
Also skip the TCK beforeEach on dialects without batching support, so the
batched journal specs pend on MySQL instead of failing CI.
Tests:
- sbt -Dpekko.persistence.r2dbc.dialect=mysql "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" - 0 failed, 74
pending
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" - 75 passed, 0
failed, PostgreSQL
- sbt "core / Test / test" - 353 passed, 0 failed, PostgreSQL
- sbt +headerCheckAll
- sbt "++3.3.8 core/Test/compile" "++3.9.0 core/Test/compile"
References: Refs #471
* fix: port useful #491 cleanups to batched write journal (#471)
Motivation:
PR #491 suggested cleanups on top of this branch. This commit adapts the
useful ones while keeping the original flush design.
Modification:
- Reuse R2dbcJournal.WriteFinished and deserializeRow instead of
duplicating them.
- Complete write requests with Done and derive the AsyncWriteJournal result
at the API boundary.
- Cancel a pending Flush timer when a flush starts so it does not flush the
next partial batch early.
- Take the batch from the queue as a Vector and drop the immutable
collection prefixes.
- Correct the JournalDao.writeEvents comment about mixed persistence ids.
- Remove the redundant batched-journal block from the MultiPluginSpec
snippet.
- Docs: paradox warning block, corrected bisection retry cost, accurate
timer and transaction wording.
Result:
Less duplicated code, a completion signal with a single meaning, and docs
that describe the actual flush semantics. The queue still flushes when
max-batch-size is reached or max-batch-time expires; the saturation re-flush on
flush completion stays as part of the capacity trigger.
Tests:
- sbt headerCheckAll
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*
org.apache.pekko.persistence.r2dbc.journal.MultiPluginSpec" - 76 passed, 0
failed, PostgreSQL 18
- scalafmt (native 3.11.5) and git diff --check
References:
Refs #471, Refs #489
* fix: address third review round on batched write journal (#471)
Motivation:
Review of b5877fa found that db_timestamp was sampled at enqueue time, a
full
queue tripped the journal circuit breaker, bisection retries ran serially
inside the flush slot, and max-batch-time defaulted below the scheduler
tick.
Modification:
Stamp rows once per flush immediately before writeEvents, publish the
timestamp
writeEvents returns, and re-stamp on bisection retries. Return queue-full as
per-message rejections in a successful Future so the circuit breaker does
not
count them, and pin the accept/reject boundary in QueueLimitSpec. Retry
bisection halves concurrently. Default max-batch-time to 10ms and document
scheduler tick rounding. Call doFlush directly on the size trigger, make
publish synchronous, and remove the stopping flag. Document the
mixed-persistence-id contract in JournalDao.writeEvents and log the actual
row
count. Document the timestamp lag bound, dialect support, failure semantics,
and the same-persistence-id db_timestamp inversion that a bisection retry
allows. Preserve the derived-from-Akka headers on the four files that copy
Akka-derived code.
Result:
Published pub-sub offsets match the stored rows, sustained queue overload no
longer opens the journal circuit breaker, and bisection retry latency is
about
one transaction per level instead of two.
Tests:
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" (76 tests, live
PostgreSQL)
- sbt headerCheckAll
- sbt docs/paradox
- scalafmt --test on the changed files
- git diff --check
- sbt +mimaReportBinaryIssues (local run skips analysis,
mimaPreviousArtifacts empty; the GitHub Check / Binary Compatibility job is
authoritative)
References:
Refs #471
* Revise license and copyright comments in R2dbcBatchJournal
Updated license information and copyright notice in R2dbcBatchJournal.scala.
* Revise license and copyright in test file
Updated license information and copyright notice in
R2dbcBatchJournalPerfManyActorsSpec.scala.
* Refactor license comments in R2dbcBatchJournalPerfSpec
Updated license information and removed redundant comments.
* Revise license and copyright comments in test file
Updated license information and copyright notice in
R2dbcBatchJournalSpec.scala.
* fix: stamp writes individually and guard stale FlushDone (#489)
Motivation:
One db_timestamp per flush could stall eventsBySlices paging, timestamps
were taken before connection acquisition, a stale FlushDone could start a
second flush after an actor restart, and max-batch-time = 0 was silently
accepted.
Modification:
Stamp each write request with its own strictly increasing microsecond
timestamp and publish each request with its stored timestamp. Carry an
incarnation generation in FlushDone and ignore messages from older generations.
Require max-batch-time > 0. Document the timestamp-to-commit lag and the
behind-current-time guidance in reference.conf and journal.md.
Result:
eventsBySlices can page through any batch, stale FlushDone cannot trigger a
concurrent flush, invalid max-batch-time fails fast, and published offsets
match stored rows.
Tests:
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" - 79 passed, 0
failed, PostgreSQL 18.4
- sbt +headerCheckAll
- sbt docs/paradox
- sbt scalafmtCheckAll scalafmtSbtCheck
- git diff --check
- sbt sortImports skipped: command not available in this project
References:
Refs #471, Refs #489
* fix: address review nits on batched write journal (#489)
Motivation:
A stale FlushDone from a previous incarnation matched no case and was
logged as an unhandled dead letter. The publish warning logged only the
exception message and lost the stack trace. The docs called the timestamp
cursor per-actor when it is JVM-wide.
Modification:
Move the FlushDone generation check into the case body so stale messages
are ignored silently. Log the publish exception as the warning cause. Reword
the scaladoc and journal.md to say the stamps are strictly increasing across
all journal actor instances in the JVM.
Result:
Stale FlushDone produces no dead-letter noise, publish failures keep their
stack trace, and the docs match the actual timestamp scope.
Tests:
- sbt "core / Test / testOnly
org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal*" - 79 passed, 0
failed, PostgreSQL 18.4
- sbt docs/paradox
- scalafmt --test on R2dbcBatchJournal.scala (3.11.5)
- git diff --check
References:
Refs #471, Refs #489
---------
Co-authored-by: PJ Fanning <[email protected]>
---
core/src/main/resources/reference.conf | 32 +-
.../persistence/r2dbc/journal/JournalDao.scala | 10 +-
.../r2dbc/journal/R2dbcBatchJournal.scala | 361 +++++++++++++++++++++
.../r2dbc/journal/BatchedJournalDialectGate.scala | 49 +++
.../journal/R2dbcBatchJournalBatchingSpec.scala | 340 +++++++++++++++++++
.../R2dbcBatchJournalFailureIsolationSpec.scala | 147 +++++++++
.../R2dbcBatchJournalPerfManyActorsSpec.scala | 68 ++++
.../r2dbc/journal/R2dbcBatchJournalPerfSpec.scala | 41 +++
.../R2dbcBatchJournalPublishTimestampSpec.scala | 141 ++++++++
.../r2dbc/journal/R2dbcBatchJournalSpec.scala | 44 +++
.../journal/R2dbcBatchJournalValidationSpec.scala | 268 +++++++++++++++
docs/src/main/paradox/journal.md | 99 ++++++
12 files changed, 1594 insertions(+), 6 deletions(-)
diff --git a/core/src/main/resources/reference.conf
b/core/src/main/resources/reference.conf
index 35c9bc1..28af9f4 100644
--- a/core/src/main/resources/reference.conf
+++ b/core/src/main/resources/reference.conf
@@ -65,6 +65,34 @@ pekko.persistence.r2dbc {
use-app-timestamp = ${pekko.persistence.r2dbc.use-app-timestamp}
}
+
+ batched-journal = ${pekko.persistence.r2dbc.journal} {
+ class = "org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+
+ db-timestamp-monotonic-increasing = on
+
+ use-app-timestamp = on
+
+ # The batched journal stamps rows before connection acquisition. Under
connection pool contention,
+ # the lag between timestamping and commit can exceed
query.behind-current-time (default 100ms).
+ # Consider raising pekko.persistence.r2dbc.query.behind-current-time
comfortably above that lag
+ # (e.g. 1s or higher) so live eventsBySlices queries do not skip events
before backtracking.
+
+ # Maximum number of write requests buffered when using the batched journal.
+ # When the queue has reached this limit, new writes are rejected per
message
+ # in the journal write results. The rejection does not count toward the
+ # journal circuit breaker.
+ max-queue-size = 10000
+
+ # Maximum number of write requests to batch when using the batched journal
+ max-batch-size = 100
+
+ # Maximum duration to wait before flushing batched writes when using the
batched journal.
+ # Must be greater than zero.
+ # Pekko timers are rounded up to whole scheduler ticks
(pekko.scheduler.tick-duration,
+ # default 10ms), so a value below the tick duration takes effect as one
tick.
+ max-batch-time = 10ms
+ }
}
// #journal-settings
@@ -319,8 +347,8 @@ pekko.persistence.r2dbc {
# move backwards when the system clock is adjusted.
db-timestamp-monotonic-increasing = off
- # Enable this for testing or workaround of
https://github.com/yugabyte/yugabyte-db/issues/10995
- # FIXME: This property will be removed when the Yugabyte issue has been
resolved.
+ # Required by the MySQL dialect and the batched journal, and a workaround of
+ # https://github.com/yugabyte/yugabyte-db/issues/10995 for Yugabyte.
use-app-timestamp = off
# Logs database calls that take longer than this duration at INFO level.
diff --git
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/JournalDao.scala
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/JournalDao.scala
index a66fb52..1d621f3 100644
---
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/JournalDao.scala
+++
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/JournalDao.scala
@@ -175,7 +175,10 @@ private[r2dbc] class JournalDao(val settings:
JournalSettings, connectionFactory
VALUES (?, ?, ?, ?, $timestampSql, ?, ?, ?, ?, ?, ?)"""
/**
- * All events must be for the same persistenceId.
+ * All events must be for the same persistenceId, except when called by
`R2dbcBatchJournal`, which mixes
+ * persistence ids in one batch. Mixing is only possible with
`db-timestamp-monotonic-increasing`, where the
+ * per-persistence-id previous sequence number is not bound, and in that
case the persistenceId is only used
+ * for logging.
*
* The returned timestamp should be the `db_timestamp` column and it is used
in published events when that feature is
* enabled.
@@ -187,7 +190,6 @@ private[r2dbc] class JournalDao(val settings:
JournalSettings, connectionFactory
def writeEvents(events: Seq[SerializedJournalRow]): Future[Instant] = {
require(events.nonEmpty)
- // it's always the same persistenceId for all events
val persistenceId = events.head.persistenceId
val previousSeqNr = events.head.seqNr - 1
@@ -252,7 +254,7 @@ private[r2dbc] class JournalDao(val settings:
JournalSettings, connectionFactory
row => row.get(0, classOf[Instant]))
if (log.isDebugEnabled())
result.foreach { _ =>
- log.debug("Wrote [{}] events for persistenceId [{}]", 1,
events.head.persistenceId)
+ log.debug("Wrote [{}] events for persistenceId [{}]", totalEvents,
events.head.persistenceId)
}
if (useTimestampFromDb) {
result
@@ -271,7 +273,7 @@ private[r2dbc] class JournalDao(val settings:
JournalSettings, connectionFactory
row => row.get(0, classOf[Instant]))
if (log.isDebugEnabled())
result.foreach { _ =>
- log.debug("Wrote [{}] events for persistenceId [{}]", 1,
events.head.persistenceId)
+ log.debug("Wrote [{}] events for persistenceId [{}]", totalEvents,
events.head.persistenceId)
}
if (useTimestampFromDb) {
result.map(_.head)(ExecutionContext.parasitic)
diff --git
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournal.scala
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournal.scala
new file mode 100644
index 0000000..9ff3247
--- /dev/null
+++
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournal.scala
@@ -0,0 +1,361 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * license agreements; and to You under the Apache License, version 2.0:
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * This file is part of the Apache Pekko project, which was derived from Akka.
+ */
+
+/*
+ * Copyright (C) 2021 - 2023 Lightbend Inc. <https://www.lightbend.com>
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import java.time.Instant
+import java.time.temporal.ChronoUnit
+
+import scala.concurrent.{ ExecutionContext, Future, Promise }
+import scala.concurrent.duration.{ Duration, FiniteDuration }
+import scala.jdk.DurationConverters.JavaDurationOps
+import scala.util.{ Failure, Success, Try }
+import scala.util.control.NonFatal
+
+import com.typesafe.config.Config
+import io.r2dbc.spi.R2dbcDataIntegrityViolationException
+import org.apache.pekko
+import pekko.Done
+import pekko.actor.Timers
+import pekko.actor.typed.ActorSystem
+import pekko.actor.typed.scaladsl.adapter._
+import pekko.annotation.InternalApi
+import pekko.event.Logging
+import pekko.persistence.AtomicWrite
+import pekko.persistence.Persistence
+import pekko.persistence.PersistentRepr
+import pekko.persistence.journal.AsyncWriteJournal
+import pekko.persistence.journal.Tagged
+import pekko.persistence.r2dbc.Dialect.{ Postgres, Yugabyte }
+import pekko.persistence.r2dbc.JournalSettings
+import pekko.persistence.r2dbc.internal.InstantFactory
+import pekko.persistence.r2dbc.internal.PubSub
+import pekko.persistence.r2dbc.journal.JournalDao.SerializedEventMetadata
+import pekko.persistence.r2dbc.journal.JournalDao.SerializedJournalRow
+import pekko.persistence.typed.PersistenceId
+import pekko.serialization.Serialization
+import pekko.serialization.SerializationExtension
+import pekko.serialization.Serializers
+import pekko.stream.scaladsl.Sink
+
+/**
+ * INTERNAL API
+ */
+@InternalApi
+private[r2dbc] object R2dbcBatchJournal {
+ private case object Flush
+ private[r2dbc] final case class FlushDone(generation: Long)
+
+ // the promise is completed with Done only after the batch containing this
request is committed;
+ // the AsyncWriteJournal result is derived from it at the API boundary
+ private final case class WriteRequest(
+ rows: Seq[SerializedJournalRow],
+ messages: Seq[AtomicWrite],
+ promise: Promise[Done]
+ )
+
+ private val generationCounter = new
java.util.concurrent.atomic.AtomicLong(0L)
+ private[r2dbc] def nextGeneration(): Long =
generationCounter.incrementAndGet()
+
+ private val lastTimestampMicros = new
java.util.concurrent.atomic.AtomicLong(0L)
+ private[r2dbc] def nextTimestamp(): Instant = {
+ val nowMicros = ChronoUnit.MICROS.between(Instant.EPOCH,
InstantFactory.now())
+ val next = lastTimestampMicros.updateAndGet(prev => math.max(prev + 1,
nowMicros))
+ Instant.EPOCH.plus(next, ChronoUnit.MICROS)
+ }
+}
+
+/**
+ * INTERNAL API
+ *
+ * Opt-in journal plugin (`pekko.persistence.r2dbc.batched-journal`) that
coalesces concurrent
+ * writes from different persistence ids into one transaction, trading up to
+ * `max-batch-time` of write latency for higher throughput at high concurrency.
+ *
+ * Mixing persistence ids in a single statement is only safe because the
plugin requires
+ * `use-app-timestamp = on` and `db-timestamp-monotonic-increasing = on`: in
that mode
+ * [[JournalDao]] does not bind the per-persistence-id previous sequence
number subselect, and
+ * timestamps come from the application clock, which therefore must not move
backwards.
+ * Batching is only supported and tested for the Postgres and Yugabyte
dialects.
+ *
+ * Writes are buffered in a bounded queue (`max-queue-size`); incoming writes
are rejected per
+ * message once the queue is full. The rejection is returned in the
per-message `Try` results,
+ * not as a failed `Future`, so a full queue does not count toward the journal
circuit breaker.
+ * Flushed batches are serialized: only one flush is in flight at a time,
which keeps
+ * same-persistence-id writes committed in order without relying on
replay-time coordination, at
+ * the cost of not using spare pool capacity. Concurrent flushing can be added
later if a single
+ * flush saturates.
+ *
+ * Each write request is stamped at flush time with the application clock
truncated to
+ * microseconds and bumped to stay strictly increasing across all journal
actor instances in the JVM.
+ * Equal `db_timestamp` values therefore never span more than one write
request within the JVM, so
+ * the `eventsBySlices` query can page through any batch regardless of its
buffer size. The
+ * stamps can lead the wall clock by at most `max-batch-size` microseconds per
flush. A single
+ * request can still contain many events when the caller uses `persistAll` or
`persistAsync`
+ * bursts, the same as the default journal.
+ *
+ * A batch that fails with a database integrity violation is retried in halves
so that only the
+ * offending persistence ids fail. Infrastructure errors fail the whole batch.
The retried
+ * halves are written concurrently and can use several pool connections at
once. Each half
+ * re-stamps its requests with new, later timestamps. A single split preserves
order, but if the
+ * half holding the earlier sequence numbers is retried after its sibling
committed, its
+ * re-stamped rows can invert the `db_timestamp` order of two
same-persistence-id writes.
+ * Replay still orders by sequence number; only timestamp-ordered read sides
see the inversion.
+ * This plugin targets many small concurrent writes.
+ */
+@InternalApi
+private[r2dbc] final class R2dbcBatchJournal(config: Config) extends
AsyncWriteJournal with Timers {
+ import R2dbcJournal.WriteFinished
+ import R2dbcJournal.deserializeRow
+ import R2dbcBatchJournal.Flush
+ import R2dbcBatchJournal.FlushDone
+ import R2dbcBatchJournal.WriteRequest
+
+ implicit val system: ActorSystem[?] = context.system.toTyped
+ implicit val ec: ExecutionContext = context.dispatcher
+
+ private val log = Logging(context.system, classOf[R2dbcBatchJournal])
+
+ private val persistenceExt = Persistence(system)
+
+ private val serialization: Serialization =
SerializationExtension(context.system)
+ private val journalSettings = JournalSettings(config)
+
+ require(journalSettings.dialect == Postgres || journalSettings.dialect ==
Yugabyte,
+ "Batching is only supported for Postgres and Yugabyte")
+ require(journalSettings.useAppTimestamp, "use-app-timestamp must be 'on'
when using R2dbcBatchJournal")
+ require(journalSettings.dbTimestampMonotonicIncreasing,
+ "db-timestamp-monotonic-increasing must be 'on' when using
R2dbcBatchJournal")
+
+ private val maxQueueSize: Int = config.getInt("max-queue-size")
+ private val maxBatchSize: Int = config.getInt("max-batch-size")
+ private val maxBatchTime: FiniteDuration =
config.getDuration("max-batch-time").toScala
+
+ require(maxQueueSize > 0, "max-queue-size must be at least 1 when using
R2dbcBatchJournal")
+ require(maxBatchSize > 0, "max-batch-size must be at least 1 when using
R2dbcBatchJournal")
+ require(maxBatchSize <= maxQueueSize, "max-batch-size must be less than or
equal to `max-queue-size`")
+ require(maxBatchTime > Duration.Zero, "max-batch-time must be greater than
zero when using R2dbcBatchJournal")
+
+ private val generation = R2dbcBatchJournal.nextGeneration()
+
+ private val journalDao = JournalDao.fromConfig(journalSettings, config)
+
+ private val pubSub: Option[PubSub] =
+ Option.when(journalSettings.journalPublishEvents)(PubSub(system))
+
+ // if there are pending writes when an actor restarts we must wait for
+ // them to complete before we can read the highest sequence number, or we
will miss it
+ private val writesInProgress = new java.util.HashMap[String, Future[?]]()
+
+ private val queue = collection.mutable.ArrayDeque[WriteRequest]()
+ private var noActiveWrite = true
+
+ private def doFlush(): Unit = {
+ // a pending timer would otherwise flush the next, partial batch early
+ timers.cancel(Flush)
+
+ val count = math.min(maxBatchSize, queue.size)
+ val writeRequests = queue.take(count).toVector
+ queue.dropInPlace(count)
+ log.debug("flushing [{}] write requests", count)
+
+ def write(requests: Vector[WriteRequest]): Future[Unit] = {
+ val stampedRequests = requests.map(request => request ->
R2dbcBatchJournal.nextTimestamp())
+
+ journalDao
+ .writeEvents(stampedRequests.flatMap {
+ case (request, timestamp) => request.rows.map(_.copy(dbTimestamp =
timestamp))
+ })
+ .map { _ =>
+ requests.foreach(_.promise.trySuccess(Done))
+ publish(stampedRequests)
+ }
+ .recoverWith {
+ case _: R2dbcDataIntegrityViolationException if requests.size > 1 =>
+ val (left, right) = requests.splitAt(requests.size / 2)
+ write(left).zipWith(write(right))((_, _) => ())
+ case exception =>
+ requests.foreach(_.promise.tryFailure(exception))
+ Future.unit
+ }
+ }
+
+ write(writeRequests).onComplete(_ => self ! FlushDone(generation))
+ }
+
+ override def receivePluginInternal: Receive = {
+ case WriteFinished(pid, f) => writesInProgress.remove(pid, f)
+ case Flush =>
+ if (noActiveWrite && queue.nonEmpty) {
+ noActiveWrite = false
+ doFlush()
+ }
+ case FlushDone(g) =>
+ if (g == generation) {
+ noActiveWrite = true
+ if (queue.size >= maxBatchSize) {
+ noActiveWrite = false
+ doFlush()
+ } else if (queue.nonEmpty && !timers.isTimerActive(Flush)) {
+ timers.startSingleTimer(Flush, Flush, maxBatchTime)
+ }
+ }
+ }
+
+ override def asyncWriteMessages(messages: Seq[AtomicWrite]):
Future[Seq[Try[Unit]]] = {
+ if (queue.length >= maxQueueSize) {
+ val queueFullFailure =
+ Failure(new IllegalStateException(s"Unable to accept the request,
max-queue-size [$maxQueueSize] reached"))
+ Future.successful(messages.map(_ => queueFullFailure))
+ } else {
+ val promise = Promise[Done]()
+
+ def atomicWrite(atomicWrite: AtomicWrite):
Try[Seq[SerializedJournalRow]] = {
+ val serialized: Try[Seq[SerializedJournalRow]] = Try {
+ atomicWrite.payload.map { pr =>
+ val (event, tags) = pr.payload match {
+ case Tagged(payload, tags) =>
+ (payload.asInstanceOf[AnyRef], tags)
+ case other =>
+ (other.asInstanceOf[AnyRef], Set.empty[String])
+ }
+
+ val entityType = PersistenceId.extractEntityType(pr.persistenceId)
+ val slice = persistenceExt.sliceForPersistenceId(pr.persistenceId)
+
+ val serialized = serialization.serialize(event).get
+ val serializer = serialization.findSerializerFor(event)
+ val manifest = Serializers.manifestFor(serializer, event)
+ val id: Int = serializer.identifier
+
+ val metadata = pr.metadata.map { meta =>
+ val m = meta.asInstanceOf[AnyRef]
+ val serializedMeta = serialization.serialize(m).get
+ val metaSerializer = serialization.findSerializerFor(m)
+ val metaManifest = Serializers.manifestFor(metaSerializer, m)
+ val id: Int = metaSerializer.identifier
+ SerializedEventMetadata(id, metaManifest, serializedMeta)
+ }
+
+ SerializedJournalRow(
+ slice,
+ entityType,
+ pr.persistenceId,
+ pr.sequenceNr,
+ JournalDao.EmptyDbTimestamp,
+ JournalDao.EmptyDbTimestamp,
+ Some(serialized),
+ id,
+ manifest,
+ pr.writerUuid,
+ tags,
+ metadata)
+ }
+ }
+
+ serialized match {
+ case Success(writes) =>
+ queue.addOne(WriteRequest(writes, Seq(atomicWrite), promise))
+
+ writesInProgress.put(writes.head.persistenceId, promise.future)
+ promise.future.onComplete { _ =>
+ self ! WriteFinished(writes.head.persistenceId, promise.future)
+ }
+
+ if (queue.size >= maxBatchSize && noActiveWrite) {
+ noActiveWrite = false
+ doFlush()
+ } else if (!timers.isTimerActive(Flush))
+ timers.startSingleTimer(Flush, Flush, maxBatchTime)
+ case Failure(exception) =>
+ promise.tryFailure(exception)
+ }
+
+ serialized
+ }
+
+ if (messages.size == 1)
+ atomicWrite(messages.head)
+ else {
+ // persistAsync case
+ // easiest to just group all into a single AtomicWrite
+ val batch = AtomicWrite(messages.flatMap(_.payload))
+ atomicWrite(batch)
+ }
+
+ // an empty result means that all messages were written, as in
R2dbcJournal
+ promise.future.map(_ => Nil)(ExecutionContext.parasitic)
+ }
+ }
+
+ private def publish(requests: Vector[(WriteRequest, Instant)]): Unit =
+ pubSub.foreach { ps =>
+ requests.foreach {
+ case (request, timestamp) =>
+ try {
+ request.messages.foreach { messages =>
+ messages.payload.foreach(pr => ps.publish(pr, timestamp))
+ }
+ } catch {
+ case NonFatal(exception) =>
+ log.warning(
+ exception,
+ "Failed to publish events for persistence id [{}]",
+ request.messages.head.persistenceId)
+ }
+ }
+ }
+
+ override def asyncDeleteMessagesTo(persistenceId: String, toSequenceNr:
Long): Future[Unit] = {
+ log.debug("asyncDeleteMessagesTo persistenceId [{}], toSequenceNr [{}]",
persistenceId, toSequenceNr)
+ journalDao.deleteMessagesTo(persistenceId, toSequenceNr)
+ }
+
+ override def asyncReplayMessages(persistenceId: String, fromSequenceNr:
Long, toSequenceNr: Long, max: Long)(
+ recoveryCallback: PersistentRepr => Unit): Future[Unit] = {
+ log.debug("asyncReplayMessages persistenceId [{}], fromSequenceNr [{}]",
persistenceId, fromSequenceNr)
+ val effectiveToSequenceNr =
+ if (max == Long.MaxValue) toSequenceNr
+ else math.min(toSequenceNr, fromSequenceNr + max - 1)
+ journalDao
+ .internalCurrentEventsByPersistenceId(persistenceId, fromSequenceNr,
effectiveToSequenceNr)
+ .runWith(Sink.foreach { row =>
+ val repr = deserializeRow(serialization, row)
+ recoveryCallback(repr)
+ })
+ .map(_ => ())
+ }
+
+ override def asyncReadHighestSequenceNr(persistenceId: String,
fromSequenceNr: Long): Future[Long] = {
+ log.debug("asyncReadHighestSequenceNr [{}] [{}]", persistenceId,
fromSequenceNr)
+ val pendingWrite = Option(writesInProgress.get(persistenceId)) match {
+ case Some(f) =>
+ log.debug("Write in progress for [{}], deferring highest seq nr until
write completed", persistenceId)
+ // we only want to make write - replay sequential, not fail if
previous write failed
+ f.recover { case _ => Done }(ExecutionContext.parasitic)
+ case None => Future.successful(Done)
+ }
+ pendingWrite.flatMap(_ => journalDao.readHighestSequenceNr(persistenceId,
fromSequenceNr))
+ }
+
+ override def postStop(): Unit = {
+ val cause = new IllegalStateException("Journal actor stopped with pending
batched writes")
+
+ queue.foreach(_.promise.tryFailure(cause))
+ writesInProgress.clear()
+ queue.clear()
+
+ super.postStop()
+ }
+
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/BatchedJournalDialectGate.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/BatchedJournalDialectGate.scala
new file mode 100644
index 0000000..a31d070
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/BatchedJournalDialectGate.scala
@@ -0,0 +1,49 @@
+/*
+ * 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.
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import org.apache.pekko
+import pekko.persistence.r2dbc.TestConfig
+import org.scalatest.BeforeAndAfterEach
+import org.scalatest.Outcome
+import org.scalatest.Pending
+import org.scalatest.TestSuite
+
+/**
+ * INTERNAL API
+ */
+private[r2dbc] trait BatchedJournalDialectGate extends TestSuite {
+
+ protected def batchedJournalDialectSupported: Boolean = {
+ val dialect =
TestConfig.config.getString("pekko.persistence.r2dbc.dialect")
+ dialect == "postgres" || dialect == "yugabyte"
+ }
+
+ override def withFixture(test: NoArgTest): Outcome =
+ if (batchedJournalDialectSupported) super.withFixture(test)
+ else Pending
+}
+
+/**
+ * INTERNAL API
+ */
+private[r2dbc] trait BatchedJournalTckDialectGate extends
BatchedJournalDialectGate with BeforeAndAfterEach {
+
+ abstract override protected def beforeEach(): Unit =
+ if (batchedJournalDialectSupported) super.beforeEach()
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalBatchingSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalBatchingSpec.scala
new file mode 100644
index 0000000..2c14392
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalBatchingSpec.scala
@@ -0,0 +1,340 @@
+/*
+ * 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.
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import scala.concurrent.duration._
+import org.apache.pekko
+import pekko.actor.testkit.typed.scaladsl.LogCapturing
+import pekko.actor.testkit.typed.scaladsl.LoggingTestKit
+import pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit
+import pekko.actor.typed.ActorRef
+import pekko.actor.typed.ActorSystem
+import pekko.actor.typed.scaladsl.adapter._
+import pekko.persistence.AtomicWrite
+import pekko.persistence.JournalProtocol.WriteMessageRejected
+import pekko.persistence.JournalProtocol.WriteMessageSuccess
+import pekko.persistence.JournalProtocol.WriteMessages
+import pekko.persistence.JournalProtocol.WriteMessagesSuccessful
+import pekko.persistence.PersistentRepr
+import pekko.persistence.r2dbc.ConnectionFactoryProvider
+import pekko.persistence.r2dbc.TestData
+import pekko.persistence.r2dbc.TestDbLifecycle
+import pekko.persistence.r2dbc.internal.R2dbcExecutor.PublisherOps
+import com.typesafe.config.Config
+import com.typesafe.config.ConfigFactory
+import org.scalatest.wordspec.AnyWordSpecLike
+
+object R2dbcBatchJournalBatchingSpec {
+
+ // writes are held until max-batch-size is reached, the long batch window
+ // guarantees that only the size trigger can complete the writes in the test
+ val sizeFlushConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc.batched-journal {
+ max-queue-size = 10
+ max-batch-size = 3
+ max-batch-time = 10s
+ }""")
+ .withFallback(R2dbcBatchJournalSpec.config)
+
+ // a single write below max-batch-size is flushed when the batch window
expires
+ val timerFlushConfig: Config = ConfigFactory
+ .parseString("pekko.persistence.r2dbc.batched-journal.max-batch-time =
500ms")
+ .withFallback(R2dbcBatchJournalSpec.config)
+
+ // the only pooled connection is held by the test so the flush started by
the first
+ // write stays in progress and the queue fills deterministically
+ val queueLimitConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.loglevel = DEBUG
+ pekko.persistence.r2dbc {
+ batched-journal {
+ max-queue-size = 3
+ max-batch-size = 1
+ max-batch-time = 500ms
+ use-connection-factory =
"pekko.persistence.r2dbc.queue-limit-test-connection-factory"
+ }
+ }
+ pekko.persistence.r2dbc.queue-limit-test-connection-factory =
${pekko.persistence.r2dbc.connection-factory} {
+ initial-size = 1
+ max-size = 1
+ acquire-timeout = 30s
+ }""")
+ .withFallback(R2dbcBatchJournalSpec.config)
+ .resolve()
+
+ // A long batch window so only the size trigger can flush, used to assert
that a stale FlushDone
+ // does not start a second flush while the first one is still in progress
+ val staleFlushDoneConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.loglevel = DEBUG
+ pekko.persistence.r2dbc {
+ batched-journal {
+ max-queue-size = 3
+ max-batch-size = 1
+ max-batch-time = 10s
+ use-connection-factory =
"pekko.persistence.r2dbc.queue-limit-test-connection-factory"
+ }
+ }
+ pekko.persistence.r2dbc.queue-limit-test-connection-factory =
${pekko.persistence.r2dbc.connection-factory} {
+ initial-size = 1
+ max-size = 1
+ acquire-timeout = 30s
+ }""")
+ .withFallback(R2dbcBatchJournalSpec.config)
+ .resolve()
+
+ def writeMessages(pid: String, seqNr: Long, event: String, replyTo:
ActorRef[Any]): WriteMessages =
+ WriteMessages(
+ Seq(AtomicWrite(PersistentRepr(event, seqNr, pid))),
+ replyTo.toClassic,
+ actorInstanceId = 1)
+}
+
+class R2dbcBatchJournalSizeFlushSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalBatchingSpec.sizeFlushConfig)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+ import R2dbcBatchJournalBatchingSpec.writeMessages
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+
+ "R2dbcBatchJournal size flush" should {
+
+ "hold writes until max-batch-size is reached and then complete them all"
in {
+ val entityType = nextEntityType()
+ val pids = (1 to 3).map(_ => nextPid(entityType))
+ val probes = pids.map(_ => createTestProbe[Any]())
+
+ // max-batch-time is 10s so the first two writes cannot complete on
their own
+ pids.take(2).zip(probes).foreach {
+ case (pid, probe) =>
+ journal ! writeMessages(pid, 1L, s"e-$pid", probe.ref)
+ }
+ probes.take(2).foreach(_.expectNoMessage(1.second))
+
+ // the third write reaches max-batch-size = 3 and triggers the flush of
the whole batch
+ journal ! writeMessages(pids(2), 1L, s"e-${pids(2)}", probes(2).ref)
+
+ pids.zip(probes).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(5.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](5.seconds).persistent.persistenceId
shouldBe pid
+ }
+ }
+
+ }
+
+}
+
+class R2dbcBatchJournalTimerFlushSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalBatchingSpec.timerFlushConfig)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+ import R2dbcBatchJournalBatchingSpec.writeMessages
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+
+ "R2dbcBatchJournal timer flush" should {
+
+ "flush a partial batch when max-batch-time expires" in {
+ val entityType = nextEntityType()
+ val pid = nextPid(entityType)
+ val probe = createTestProbe[Any]()
+
+ journal ! writeMessages(pid, 1L, s"e-$pid", probe.ref)
+
+ // max-batch-size is 100 so the write cannot complete before the 500ms
batch window expires
+ probe.expectNoMessage(200.millis)
+ probe.expectMessage(5.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](5.seconds).persistent.persistenceId
shouldBe pid
+ }
+
+ }
+
+}
+
+class R2dbcBatchJournalStaleFlushDoneSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalBatchingSpec.staleFlushDoneConfig)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+ import R2dbcBatchJournalBatchingSpec.writeMessages
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+ private val journalConnectionFactory =
+ ConnectionFactoryProvider(system).connectionFactoryFor(
+ "pekko.persistence.r2dbc.queue-limit-test-connection-factory")
+
+ "R2dbcBatchJournal stale FlushDone" should {
+
+ "ignore a FlushDone from a previous incarnation" in {
+ val entityType = nextEntityType()
+ val pids = (1 to 2).map(_ => nextPid(entityType))
+ val probes = pids.map(_ => createTestProbe[Any]())
+
+ // hold the only pooled connection so the flush of the first write stays
in progress
+ val blockingConnection =
journalConnectionFactory.create().asFuture().futureValue
+
+ // the first write starts a flush that is blocked waiting for the held
connection
+ LoggingTestKit.debug("flushing [1] write requests").expect {
+ journal ! writeMessages(pids(0), 1L, s"e-${pids(0)}", probes(0).ref)
+ }
+
+ // a FlushDone carrying a generation from a previous incarnation must be
ignored:
+ // no second flush may start while the first one is still in progress
+ journal ! R2dbcBatchJournal.FlushDone(-1L)
+ LoggingTestKit.debug("flushing").withOccurrences(0).expect {
+ journal ! writeMessages(pids(1), 1L, s"e-${pids(1)}", probes(1).ref)
+ }
+
+ // release the connection so the blocked flush and then the second write
can complete
+ blockingConnection.close().asFuture().futureValue
+
+ pids.zip(probes).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pid
+ }
+ }
+
+ }
+
+}
+
+class R2dbcBatchJournalQueueLimitSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalBatchingSpec.queueLimitConfig)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+ import R2dbcBatchJournalBatchingSpec.writeMessages
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+ private val journalConnectionFactory =
+ ConnectionFactoryProvider(system).connectionFactoryFor(
+ "pekko.persistence.r2dbc.queue-limit-test-connection-factory")
+
+ "R2dbcBatchJournal queue limit" should {
+
+ "reject writes beyond the queue limit" in {
+ val entityType = nextEntityType()
+ val pids = (1 to 6).map(_ => nextPid(entityType))
+ val probes = pids.map(_ => createTestProbe[Any]())
+
+ // hold the only pooled connection so the flush of the first write stays
in progress
+ val blockingConnection =
journalConnectionFactory.create().asFuture().futureValue
+
+ // wait until the flush has dequeued the first write. The queue then has
room for
+ // exactly max-queue-size (3) more writes, so writes 2 - 4 are accepted
and
+ // writes 5 and 6 are rejected
+ LoggingTestKit.debug("flushing [1] write requests").expect {
+ journal ! writeMessages(pids(0), 1L, s"e-${pids(0)}", probes(0).ref)
+ }
+
+ pids.drop(1).zip(probes.drop(1)).foreach {
+ case (pid, probe) =>
+ journal ! writeMessages(pid, 1L, s"e-$pid", probe.ref)
+ }
+
+ // release the connection so the accepted writes can complete
+ blockingConnection.close().asFuture().futureValue
+
+ pids.take(4).zip(probes.take(4)).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pid
+ }
+
+ pids.takeRight(2).zip(probes.takeRight(2)).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+ val rejected =
probe.expectMessageType[WriteMessageRejected](10.seconds)
+ rejected.message.persistenceId shouldBe pid
+ rejected.cause.getMessage shouldBe "Unable to accept the request,
max-queue-size [3] reached"
+ }
+ }
+
+ "not trip the journal circuit breaker when the queue is full" in {
+ val entityType = nextEntityType()
+ val pids = (1 to 17).map(_ => nextPid(entityType))
+ val probes = pids.map(_ => createTestProbe[Any]())
+
+ // hold the only pooled connection so the flush of the first write stays
in progress
+ val blockingConnection =
journalConnectionFactory.create().asFuture().futureValue
+
+ // wait until the flush has dequeued the first write
+ LoggingTestKit.debug("flushing [1] write requests").expect {
+ journal ! writeMessages(pids(0), 1L, s"e-${pids(0)}", probes(0).ref)
+ }
+
+ // writes 2 - 4 fill the queue; writes 5 - 16 are rejected, which is
more than the
+ // default circuit-breaker max-failures (10). The rejections are
returned as
+ // per-message rejections in a successful Future, so they must not count
as
+ // circuit-breaker failures. The rejections are delivered immediately,
but with
+ // write-response-global-order = on (the default) the AsyncWriteJournal
resequencer
+ // orders responses by request arrival, so they only reach the probes
after the
+ // blocked write 1 has completed.
+ pids.slice(1, 16).zip(probes.slice(1, 16)).foreach {
+ case (pid, probe) =>
+ journal ! writeMessages(pid, 1L, s"e-$pid", probe.ref)
+ }
+
+ // release the connection so the accepted writes complete
+ blockingConnection.close().asFuture().futureValue
+
+ pids.take(4).zip(probes.take(4)).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pid
+ }
+
+ pids.slice(4, 16).zip(probes.slice(4, 16)).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+ val rejected =
probe.expectMessageType[WriteMessageRejected](10.seconds)
+ rejected.message.persistenceId shouldBe pid
+ rejected.cause.getMessage shouldBe "Unable to accept the request,
max-queue-size [3] reached"
+ }
+
+ // if the rejections had tripped the circuit breaker, this write would
fail
+ journal ! writeMessages(pids(16), 1L, s"e-${pids(16)}", probes(16).ref)
+ probes(16).expectMessage(10.seconds, WriteMessagesSuccessful)
+
probes(16).expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pids(16)
+ }
+
+ }
+
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalFailureIsolationSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalFailureIsolationSpec.scala
new file mode 100644
index 0000000..60f23bb
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalFailureIsolationSpec.scala
@@ -0,0 +1,147 @@
+/*
+ * 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.
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import java.time.Instant
+import scala.concurrent.duration._
+import org.apache.pekko
+import pekko.actor.testkit.typed.scaladsl.LogCapturing
+import pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit
+import pekko.actor.testkit.typed.scaladsl.TestProbe
+import pekko.actor.typed.ActorRef
+import pekko.actor.typed.ActorSystem
+import pekko.actor.typed.scaladsl.adapter._
+import pekko.persistence.AtomicWrite
+import pekko.persistence.JournalProtocol.WriteMessageFailure
+import pekko.persistence.JournalProtocol.WriteMessageSuccess
+import pekko.persistence.JournalProtocol.WriteMessages
+import pekko.persistence.JournalProtocol.WriteMessagesFailed
+import pekko.persistence.JournalProtocol.WriteMessagesSuccessful
+import pekko.persistence.PersistentRepr
+import pekko.persistence.r2dbc.TestData
+import pekko.persistence.r2dbc.TestDbLifecycle
+import pekko.persistence.r2dbc.internal.PayloadCodec
+import pekko.persistence.r2dbc.internal.PayloadCodec.RichRow
+import pekko.serialization.SerializationExtension
+import com.typesafe.config.Config
+import com.typesafe.config.ConfigFactory
+import io.r2dbc.spi.R2dbcDataIntegrityViolationException
+import org.scalatest.wordspec.AnyWordSpecLike
+
+object R2dbcBatchJournalFailureIsolationSpec {
+
+ // long batch window so that the concurrent writes in the test are guaranteed
+ // to be coalesced into one batch flush
+ val config: Config = ConfigFactory
+ .parseString("pekko.persistence.r2dbc.batched-journal.max-batch-time = 1s")
+ .withFallback(R2dbcBatchJournalSpec.config)
+
+ final case class StoredRow(pid: String, seqNr: Long, event: String,
dbTimestamp: Instant)
+}
+
+class R2dbcBatchJournalFailureIsolationSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalFailureIsolationSpec.config)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private implicit val journalPayloadCodec: PayloadCodec =
journalSettings.journalPayloadCodec
+ private val serialization = SerializationExtension(system)
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+ import R2dbcBatchJournalFailureIsolationSpec.StoredRow
+
+ private def sendWrite(pid: String, seqNr: Long, event: String, replyTo:
ActorRef[Any]): Unit =
+ journal ! WriteMessages(
+ Seq(AtomicWrite(PersistentRepr(event, seqNr, pid))),
+ replyTo.toClassic,
+ actorInstanceId = 1)
+
+ private def expectSuccess(probe: TestProbe[Any], pid: String): Unit = {
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pid
+ }
+
+ private def storedRows(): IndexedSeq[StoredRow] =
+ r2dbcExecutor
+ .select[StoredRow]("test")(
+ connection =>
+ connection.createStatement(
+ s"select persistence_id, seq_nr, event_ser_id, event_ser_manifest,
event_payload, db_timestamp " +
+ s"from ${journalSettings.journalTableWithSchema}"),
+ row => {
+ val event = serialization
+ .deserialize(
+ row.getPayload("event_payload"),
+ row.get[Integer]("event_ser_id", classOf[Integer]),
+ row.get("event_ser_manifest", classOf[String]))
+ .get
+ .asInstanceOf[String]
+ StoredRow(
+ row.get("persistence_id", classOf[String]),
+ row.get[java.lang.Long]("seq_nr",
classOf[java.lang.Long]).longValue(),
+ event,
+ row.get("db_timestamp", classOf[Instant]))
+ })
+ .futureValue
+
+ "R2dbcBatchJournal failure isolation" should {
+
+ "fail only the persistence id violating the unique constraint when batched
with other writes" in {
+ val entityType = nextEntityType()
+ val pidA = nextPid(entityType)
+ val pidB = nextPid(entityType)
+ val pidC = nextPid(entityType)
+
+ val probeA = createTestProbe[Any]()
+ val probeB = createTestProbe[Any]()
+ val probeC = createTestProbe[Any]()
+
+ // seed pidA seqNr 1 so that the duplicate write below violates PRIMARY
KEY(persistence_id, seq_nr)
+ sendWrite(pidA, 1L, "a1", probeA.ref)
+ expectSuccess(probeA, pidA)
+
+ // these three writes arrive within the 1 second batch window and are
flushed as one batch,
+ // the duplicate for pidA makes the batch fail with a unique constraint
violation,
+ // bisection retries the halves so that only pidA fails
+ sendWrite(pidA, 1L, "a1-duplicate", probeA.ref)
+ sendWrite(pidB, 1L, "b1", probeB.ref)
+ sendWrite(pidC, 1L, "c1", probeC.ref)
+
+ val failed = probeA.expectMessageType[WriteMessagesFailed](10.seconds)
+ failed.cause shouldBe a[R2dbcDataIntegrityViolationException]
+ val failure = probeA.expectMessageType[WriteMessageFailure](10.seconds)
+ failure.cause shouldBe a[R2dbcDataIntegrityViolationException]
+ failure.message.persistenceId shouldBe pidA
+ failure.message.sequenceNr shouldBe 1L
+
+ expectSuccess(probeB, pidB)
+ expectSuccess(probeC, pidC)
+
+ val rows = storedRows()
+ rows.filter(_.pid == pidA).map(r => (r.seqNr, r.event)) shouldBe
Vector((1L, "a1"))
+ rows.filter(_.pid == pidB).map(r => (r.seqNr, r.event)) shouldBe
Vector((1L, "b1"))
+ rows.filter(_.pid == pidC).map(r => (r.seqNr, r.event)) shouldBe
Vector((1L, "c1"))
+ }
+
+ }
+
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfManyActorsSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfManyActorsSpec.scala
new file mode 100644
index 0000000..7518ddb
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfManyActorsSpec.scala
@@ -0,0 +1,68 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * license agreements; and to You under the Apache License, version 2.0:
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * This file is part of the Apache Pekko project, which was derived from Akka.
+ */
+
+/*
+ * Copyright (C) 2021 - 2023 Lightbend Inc. <https://www.lightbend.com>
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import scala.concurrent.duration._
+
+import org.apache.pekko
+import pekko.actor.Props
+import pekko.actor.typed.ActorSystem
+import pekko.actor.typed.scaladsl.adapter._
+import pekko.persistence.CapabilityFlag
+import pekko.persistence.journal.JournalPerfSpec
+import pekko.persistence.journal.JournalPerfSpec.BenchActor
+import pekko.persistence.journal.JournalPerfSpec.Cmd
+import pekko.persistence.journal.JournalPerfSpec.ResetCounter
+import pekko.persistence.r2dbc.TestDbLifecycle
+import pekko.testkit.TestProbe
+
+class R2dbcBatchJournalPerfManyActorsSpec extends
JournalPerfSpec(R2dbcBatchJournalPerfSpec.config)
+ with TestDbLifecycle
+ with BatchedJournalTckDialectGate {
+ override def eventsCount: Int = 10
+
+ override def measurementIterations: Int = 2 // increase when testing for real
+
+ override def awaitDurationMillis: Long = 60.seconds.toMillis
+
+ override protected def supportsRejectingNonSerializableObjects:
CapabilityFlag = CapabilityFlag.off()
+
+ override def typedSystem: ActorSystem[?] = system.toTyped
+
+ def actorCount = 20 // increase when testing for real
+
+ private val commands = Vector(1 to eventsCount: _*)
+
+ "A PersistentActor's performance with journal batching" must {
+ s"measure: persist()-ing $eventsCount events for $actorCount actors" in {
+ val testProbe = TestProbe()
+ val replyAfter = eventsCount
+ def createBenchActor(actorNumber: Int) =
+ system.actorOf(Props(classOf[BenchActor], s"$pid-$actorNumber",
testProbe.ref, replyAfter))
+ val actors = 1.to(actorCount).map(createBenchActor)
+
+ measure(d => s"Persist()-ing $eventsCount * $actorCount took
${d.toMillis} ms") {
+ for (cmd <- commands; actor <- actors) {
+ actor ! Cmd("p", cmd)
+ }
+ for (_ <- actors) {
+ testProbe.expectMsg(awaitDurationMillis.millis, commands.last)
+ }
+ for (actor <- actors) {
+ actor ! ResetCounter
+ }
+ }
+ }
+ }
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfSpec.scala
new file mode 100644
index 0000000..39413d5
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPerfSpec.scala
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * license agreements; and to You under the Apache License, version 2.0:
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * This file is part of the Apache Pekko project, which was derived from Akka.
+ */
+
+/*
+ * Copyright (C) 2021 - 2023 Lightbend Inc. <https://www.lightbend.com>
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import com.typesafe.config.Config
+
+import org.apache.pekko.actor.typed.ActorSystem
+import org.apache.pekko.actor.typed.scaladsl.adapter.ClassicActorSystemOps
+import org.apache.pekko.persistence.CapabilityFlag
+import org.apache.pekko.persistence.journal.JournalPerfSpec
+import org.apache.pekko.persistence.r2dbc.TestDbLifecycle
+
+import scala.concurrent.duration.DurationInt
+
+object R2dbcBatchJournalPerfSpec {
+ val config: Config = R2dbcBatchJournalSpec.config
+}
+
+class R2dbcBatchJournalPerfSpec extends
JournalPerfSpec(R2dbcBatchJournalPerfSpec.config) with TestDbLifecycle
+ with BatchedJournalTckDialectGate {
+ override def eventsCount: Int = 200
+
+ override def measurementIterations: Int = 2 // increase when testing for real
+
+ override def awaitDurationMillis: Long = 60.seconds.toMillis
+
+ override protected def supportsRejectingNonSerializableObjects:
CapabilityFlag = CapabilityFlag.off()
+
+ override def typedSystem: ActorSystem[?] = system.toTyped
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPublishTimestampSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPublishTimestampSpec.scala
new file mode 100644
index 0000000..df4e7f2
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalPublishTimestampSpec.scala
@@ -0,0 +1,141 @@
+/*
+ * 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.
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import java.time.Instant
+import scala.concurrent.ExecutionContext
+import scala.concurrent.duration._
+import org.apache.pekko
+import pekko.actor.testkit.typed.scaladsl.LogCapturing
+import pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit
+import pekko.actor.typed.ActorRef
+import pekko.actor.typed.ActorSystem
+import pekko.actor.typed.internal.pubsub.TopicImpl
+import pekko.actor.typed.pubsub.Topic
+import pekko.actor.typed.scaladsl.adapter._
+import pekko.persistence.AtomicWrite
+import pekko.persistence.JournalProtocol.WriteMessageSuccess
+import pekko.persistence.JournalProtocol.WriteMessages
+import pekko.persistence.JournalProtocol.WriteMessagesSuccessful
+import pekko.persistence.PersistentRepr
+import pekko.persistence.query.TimestampOffset
+import pekko.persistence.query.typed.EventEnvelope
+import pekko.persistence.r2dbc.TestData
+import pekko.persistence.r2dbc.TestDbLifecycle
+import pekko.persistence.r2dbc.internal.EnvelopeOrigin
+import pekko.persistence.r2dbc.internal.PubSub
+import com.typesafe.config.Config
+import com.typesafe.config.ConfigFactory
+import org.scalatest.wordspec.AnyWordSpecLike
+
+object R2dbcBatchJournalPublishTimestampSpec {
+
+ // long batch window so that the staggered writes below are guaranteed to be
+ // coalesced into one batch flush
+ val config: Config = ConfigFactory
+ .parseString("pekko.persistence.r2dbc.batched-journal.max-batch-time = 1s")
+ .withFallback(R2dbcBatchJournalSpec.config)
+}
+
+class R2dbcBatchJournalPublishTimestampSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalPublishTimestampSpec.config)
+ with AnyWordSpecLike
+ with TestDbLifecycle
+ with TestData
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ override def typedSystem: ActorSystem[?] = system
+
+ private implicit val ec: ExecutionContext = system.executionContext
+
+ private lazy val journal =
persistenceExt.journalFor("pekko.persistence.r2dbc.batched-journal")
+
+ private def writeMessages(pid: String, seqNr: Long, event: String, replyTo:
ActorRef[Any]): WriteMessages =
+ WriteMessages(
+ Seq(AtomicWrite(PersistentRepr(event, seqNr, pid))),
+ replyTo.toClassic,
+ actorInstanceId = 1)
+
+ private def storedTimestamps(): Map[String, Instant] =
+ r2dbcExecutor
+ .select[(String, Instant)]("test")(
+ connection =>
+ connection.createStatement(
+ s"select persistence_id, db_timestamp from
${journalSettings.journalTableWithSchema}"),
+ row => row.get("persistence_id", classOf[String]) ->
row.get("db_timestamp", classOf[Instant]))
+ .futureValue
+ .toMap
+
+ "R2dbcBatchJournal publish" should {
+
+ "publish each write of a coalesced batch with its own stored db timestamp"
in {
+ val entityType = nextEntityType()
+ val pids = (1 to 3).map(_ => nextPid(entityType))
+
+ val envelopeProbe = createTestProbe[EventEnvelope[String]]()
+ val topics = pids
+ .map(pid => PubSub(system).eventTopic[String](entityType,
persistenceExt.sliceForPersistenceId(pid)))
+ .toSet
+ topics.foreach(_ ! Topic.Subscribe(envelopeProbe.ref))
+
+ // wait until the subscriptions are established
+ val statsProbe = createTestProbe[TopicImpl.TopicStats]()
+ topics.foreach { topic =>
+ eventually {
+ topic ! TopicImpl.GetTopicStats(statsProbe.ref)
+ statsProbe.receiveMessage().localSubscriberCount shouldBe 1
+ }
+ }
+
+ // stagger the writes slightly; they all fall inside the 1 second batch
window
+ // and are coalesced into one flush
+ val probes = pids.map(_ => createTestProbe[Any]())
+ pids.zip(probes).zipWithIndex.foreach {
+ case ((pid, probe), i) =>
+ system.scheduler.scheduleOnce((i * 20).millis,
+ () => journal ! writeMessages(pid, 1L, s"e-$i", probe.ref))
+ }
+
+ pids.zip(probes).foreach {
+ case (pid, probe) =>
+ probe.expectMessage(10.seconds, WriteMessagesSuccessful)
+
probe.expectMessageType[WriteMessageSuccess](10.seconds).persistent.persistenceId
shouldBe pid
+ }
+
+ val stored = storedTimestamps()
+ stored.keySet shouldBe pids.toSet
+ // Flush stamps each coalesced write with a strictly increasing stored
timestamp, which its published envelope carries.
+ stored(pids(0)).isBefore(stored(pids(1))) shouldBe true
+ stored(pids(1)).isBefore(stored(pids(2))) shouldBe true
+
+ val envelopes = envelopeProbe.receiveMessages(3, 10.seconds)
+ envelopes.map(_.persistenceId).toSet shouldBe pids.toSet
+ envelopes.foreach { env =>
+ withClue(s"pid [${env.persistenceId}]: ") {
+ env.source shouldBe EnvelopeOrigin.SourcePubSub
+ env.sequenceNr shouldBe 1L
+ env.event shouldBe s"e-${pids.indexOf(env.persistenceId)}"
+ env.offset.asInstanceOf[TimestampOffset].timestamp shouldBe
stored(env.persistenceId)
+ }
+ }
+ }
+
+ }
+
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalSpec.scala
new file mode 100644
index 0000000..0b9adf3
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalSpec.scala
@@ -0,0 +1,44 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * license agreements; and to You under the Apache License, version 2.0:
+ *
+ * https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * This file is part of the Apache Pekko project, which was derived from Akka.
+ */
+
+/*
+ * Copyright (C) 2021 - 2023 Lightbend Inc. <https://www.lightbend.com>
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import com.typesafe.config.{ Config, ConfigFactory }
+import org.apache.pekko.actor.typed.ActorSystem
+import org.apache.pekko.actor.typed.scaladsl.adapter._
+import org.apache.pekko.persistence.CapabilityFlag
+import org.apache.pekko.persistence.journal.JournalSpec
+import org.apache.pekko.persistence.r2dbc.TestDbLifecycle
+
+object R2dbcBatchJournalSpec {
+ val config: Config = ConfigFactory.parseString(
+ """
+ |pekko.persistence.r2dbc {
+ | use-app-timestamp = on
+ | db-timestamp-monotonic-increasing = on
+ | batched-journal {
+ | class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ | use-app-timestamp = on
+ | db-timestamp-monotonic-increasing = on
+ | }
+ |}
+ |pekko.persistence.journal.plugin =
"pekko.persistence.r2dbc.batched-journal"
+ |""".stripMargin
+ ).withFallback(R2dbcJournalSpec.config)
+}
+
+class R2dbcBatchJournalSpec extends JournalSpec(R2dbcBatchJournalSpec.config)
with TestDbLifecycle
+ with BatchedJournalTckDialectGate {
+ override protected def supportsRejectingNonSerializableObjects:
CapabilityFlag = CapabilityFlag.off()
+ override def typedSystem: ActorSystem[?] = system.toTyped
+}
diff --git
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalValidationSpec.scala
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalValidationSpec.scala
new file mode 100644
index 0000000..7e0ce58
--- /dev/null
+++
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/journal/R2dbcBatchJournalValidationSpec.scala
@@ -0,0 +1,268 @@
+/*
+ * 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.
+ */
+
+package org.apache.pekko.persistence.r2dbc.journal
+
+import org.apache.pekko
+import pekko.actor.testkit.typed.scaladsl.LogCapturing
+import pekko.actor.testkit.typed.scaladsl.LoggingTestKit
+import pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit
+import pekko.persistence.Persistence
+import pekko.persistence.r2dbc.TestConfig
+import com.typesafe.config.Config
+import com.typesafe.config.ConfigFactory
+import org.scalatest.wordspec.AnyWordSpecLike
+
+object R2dbcBatchJournalValidationSpec {
+
+ // TestConfig.config is resolved, so its journal block has the reference.conf
+ // substitutions frozen; the journal level settings must be overridden
explicitly
+ val zeroBatchSizeConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ max-batch-size = 0
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val appTimestampOffConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = off
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = off
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val zeroQueueSizeConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ max-queue-size = 0
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val batchSizeExceedsQueueSizeConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ max-queue-size = 1
+ max-batch-size = 2
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val zeroBatchTimeConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ max-batch-time = 0ms
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val negativeBatchTimeConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ max-batch-time = -1ms
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val monotonicIncreasingOffConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ db-timestamp-monotonic-increasing = off
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = off
+ }
+ }""")
+ .withFallback(TestConfig.config)
+
+ val mysqlDialectConfig: Config = ConfigFactory
+ .parseString("""
+ pekko.persistence.r2dbc {
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ batched-journal {
+ class =
"org.apache.pekko.persistence.r2dbc.journal.R2dbcBatchJournal"
+ use-app-timestamp = on
+ db-timestamp-monotonic-increasing = on
+ dialect = mysql
+ }
+ }""")
+ .withFallback(TestConfig.config)
+}
+
+class R2dbcBatchJournalZeroBatchSizeSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.zeroBatchSizeConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when max-batch-size is less than 1" in {
+ LoggingTestKit.error("max-batch-size must be at least 1").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalZeroQueueSizeSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.zeroQueueSizeConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when max-queue-size is less than 1" in {
+ LoggingTestKit.error("max-queue-size must be at least 1").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalBatchSizeExceedsQueueSizeSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.batchSizeExceedsQueueSizeConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when max-batch-size exceeds max-queue-size" in {
+ LoggingTestKit.error("max-batch-size must be less than or equal to
`max-queue-size`").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalZeroBatchTimeSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.zeroBatchTimeConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when max-batch-time is zero" in {
+ LoggingTestKit.error("max-batch-time must be greater than zero").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalNegativeBatchTimeSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.negativeBatchTimeConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when max-batch-time is negative" in {
+ LoggingTestKit.error("max-batch-time must be greater than zero").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalAppTimestampOffSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.appTimestampOffConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when use-app-timestamp is off" in {
+ LoggingTestKit.error("use-app-timestamp must be 'on'").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalMonotonicIncreasingOffSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.monotonicIncreasingOffConfig)
+ with AnyWordSpecLike
+ with LogCapturing
+ with BatchedJournalDialectGate {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when db-timestamp-monotonic-increasing is off" in {
+ LoggingTestKit.error("db-timestamp-monotonic-increasing must be
'on'").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
+
+class R2dbcBatchJournalMysqlDialectSpec
+ extends
ScalaTestWithActorTestKit(R2dbcBatchJournalValidationSpec.mysqlDialectConfig)
+ with AnyWordSpecLike
+ with LogCapturing {
+
+ "R2dbcBatchJournal validation" should {
+
+ "fail fast when the dialect does not support batching" in {
+ LoggingTestKit.error("Batching is only supported for Postgres and
Yugabyte").expect {
+
Persistence(system).journalFor("pekko.persistence.r2dbc.batched-journal")
+ }
+ }
+ }
+}
diff --git a/docs/src/main/paradox/journal.md b/docs/src/main/paradox/journal.md
index c224a07..5209f3a 100644
--- a/docs/src/main/paradox/journal.md
+++ b/docs/src/main/paradox/journal.md
@@ -31,6 +31,105 @@ The following can be overridden in your `application.conf`
for the journal speci
@@snip [reference.conf](/core/src/main/resources/reference.conf)
{#journal-settings}
+## Batched Journal
+
+@@@ warning { title="Experimental" }
+
+This feature is experimental and not recommended for production unless it has
been thoroughly road tested by the
+user in their own test environments.
+
+@@@
+
+The default journal writes each incoming write request with its own statement
and commit. The batched journal
+plugin (`R2dbcBatchJournal`) instead coalesces concurrent write requests from
different persistence ids into one
+transaction: the rows are written as a batch of one cached prepared statement
and committed once. This reduces
+commits and statement prepares when many persistence ids write small events at
the same time. It adds latency and
+changes failure behavior, see @ref:[Tradeoffs](#tradeoffs).
+
+The batched journal requires `use-app-timestamp` and
`db-timestamp-monotonic-increasing`, which the
+`batched-journal` configuration block enables for this plugin. This is the
same timestamp mode that the
+MySQL dialect requires. With `db-timestamp-monotonic-increasing` the database
does not enforce increasing
+timestamps per persistence id, so the application clock must not move
backwards between two writes of the
+same entity. The backtracking queries of @ref:[eventsBySlices](query.md)
recover events that were stored
+with an out-of-order timestamp. Each write request is stamped when the batch
is flushed, just before the
+insert, with the application clock truncated to microseconds and bumped to
stay strictly increasing across
+all journal actor instances in the JVM. Equal `db_timestamp` values therefore
never span more than one write request within the
+JVM (a single request can still contain several events with `persistAll` or
`persistAsync`), so
+the `eventsBySlices` query can page through any batch regardless of its buffer
size. The stamps can lead
+the wall clock by at most `max-batch-size` microseconds per flush. The lag
between the timestamp and the
+commit is bounded by the connection acquisition plus one transaction. Keep
`query.behind-current-time`
+comfortably above that lag. Batching is only supported and tested for the
Postgres and Yugabyte dialects.
+
+### Batched Journal Configuration
+
+To enable the batched journal, point the journal plugin at the
`batched-journal` block and update
+`application.conf`:
+
+```
+pekko.persistence.journal.plugin = "pekko.persistence.r2dbc.batched-journal"
+
+pekko.persistence.r2dbc.batched-journal {
+ max-queue-size = 10000 # optional, default value
+ max-batch-size = 100 # optional, default value
+ max-batch-time = 10ms # optional, default value
+}
+
+# The lag between timestamp and commit includes connection acquisition.
+# Under pool contention, raise query behind-current-time above the expected
lag.
+# pekko.persistence.r2dbc.query.behind-current-time = 1s
+```
+
+The batched journal uses the following settings, in addition to the settings
of the default journal:
+
+- `max-queue-size`: Maximum number of write requests buffered before they are
flushed. When the queue has
+ reached this limit, new writes are rejected per message in the journal write
results, without failing the
+ request `Future`, so the rejection does not count toward the journal circuit
breaker. This is the standard
+ journal rejection handling: a classic persistent actor by default logs the
rejection in `onPersistRejected`
+ and continues without the event being stored, and a typed persistent actor
restarts with an
+ `EventRejectedException`. Must be at least 1. `max-batch-size` must be less
than or equal to
+ `max-queue-size`.
+- `max-batch-size`: Maximum number of write requests in one batch. One request
can contain several events when
+ the persistent actor uses `persistAll` or `persistAsync`. A batch is flushed
when this many requests are
+ buffered.
+- `max-batch-time`: Maximum time a write request is buffered. If the batch
does not reach `max-batch-size`
+ first, it is flushed when this duration has elapsed, even if the batch holds
only one request. The default
+ is 10ms. Must be greater than zero. Pekko timers are rounded up to whole
scheduler ticks
+ (`pekko.scheduler.tick-duration`, default 10ms), so a value below the tick
duration takes effect as one
+ tick. Lowering `tick-duration` changes the timer resolution for the entire
actor system.
+
+### Tradeoffs
+
+Latency:
+A write completes when its batch is flushed. A batch is flushed when
`max-batch-size` requests are buffered,
+immediately, or when `max-batch-time` has elapsed, even if the batch holds
only one request. The maximum time
+a request is buffered is therefore the duration of one flush plus
`max-batch-time`. When `max-batch-size`
+requests are already buffered when a flush completes, the next flush starts
immediately, so batches form back
+to back.
+
+Failures:
+Writes of different persistence ids share one database statement. If the
database rejects a statement because of
+a single persistence id, for example a duplicate sequence number caused by a
zombie writer, the batched journal
+retries the batch in halves until only the offending write fails. The retry
halves run concurrently and can
+briefly use several pool connections at once. The other persistent actors wait
while the retries run. Failures
+that are not caused by a single persistence id, for example a lost database
connection, fail all writes of the
+batch. The affected persistent actors see a journal write failure, as with the
default journal: a classic
+persistent actor stops, a typed persistent actor restarts. Isolating a single
offending write costs about
+2 * log2(`max-batch-size`) additional statements; only when many writes in the
batch are offending does the
+retry approach twice `max-batch-size` statements, which is the number of
statements the default journal would
+have used for the same writes.
+
+Memory:
+`max-batch-size` limits the number of requests in one batch, not the number of
events, and the queue is limited
+by `max-queue-size`, which rejects writes once the limit is reached. A single
request can contain an arbitrary
+number of events when the persistent actor uses `persistAll` or
`persistAsync`; neither journal caps that, as in
+the default journal. With `persist()` each persistent actor has at most one
outstanding write, so the queue grows
+with the number of actively writing actors. Buffered writes are held in memory
until they are flushed.
+
+When to use:
+Batching is most effective when many persistence ids concurrently write small
events. With a low number of
+concurrent writers, or with large events, the default journal performs better
because it adds no buffering
+delay.
+
## Deletes
The journal supports deletes through hard deletes, which means the journal
entries are actually deleted from the database.
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]