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]

Reply via email to