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 5d02110  fix: backtracking queries and filtered events (#487)
5d02110 is described below

commit 5d021104677ce230dff385610b11b8f316ae8197
Author: PJ Fanning <[email protected]>
AuthorDate: Sun Sep 6 12:30:22 2026 +0100

    fix: backtracking queries and filtered events (#487)
    
    Port of akka/akka-persistence-r2dbc#445.
    
    Backtracking queries can stop too early, creating a gap between the normal
    db queries and the backtracking queries that is longer than the backtracking
    window. The switch-from-backtracking check compared the row count of the 
past
    query to bufferSize - 1, where the 1 accounts for the event at the initial
    timestamp that has already been seen and is filtered from the result. When
    events share a timestamp, more than one event is filtered out, so the query
    switched out of backtracking before catching up.
    
    Track the seen count for the latest timestamp while backtracking and carry 
it
    over as the number of events expected to be filtered by the next query.
    
    Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
 .../persistence/r2dbc/internal/BySliceQuery.scala  | 30 ++++++++++--
 .../query/EventsBySliceBacktrackingSpec.scala      | 53 +++++++++++++++++++++-
 2 files changed, 77 insertions(+), 6 deletions(-)

diff --git 
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/BySliceQuery.scala
 
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/BySliceQuery.scala
index 893adaf..d447c1f 100644
--- 
a/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/BySliceQuery.scala
+++ 
b/core/src/main/scala/org/apache/pekko/persistence/r2dbc/internal/BySliceQuery.scala
@@ -45,7 +45,7 @@ import org.slf4j.Logger
 
   object QueryState {
     val empty: QueryState =
-      QueryState(TimestampOffset.Zero, 0, 0, 0, 0, backtrackingCount = 0, 
TimestampOffset.Zero, Buckets.empty)
+      QueryState(TimestampOffset.Zero, 0, 0, 0, 0, backtrackingCount = 0, 
TimestampOffset.Zero, 0, 0, Buckets.empty)
   }
 
   final case class QueryState(
@@ -56,6 +56,8 @@ import org.slf4j.Logger
       idleCount: Long,
       backtrackingCount: Int,
       latestBacktracking: TimestampOffset,
+      latestBacktrackingSeenCount: Int,
+      backtrackingExpectFiltered: Int,
       buckets: Buckets) {
 
     def backtracking: Boolean = backtrackingCount > 0
@@ -316,11 +318,27 @@ import org.slf4j.Logger
           throw new IllegalArgumentException(
             s"Unexpected offset [$offset] before latestBacktracking 
[${state.latestBacktracking}].")
 
-        state.copy(latestBacktracking = offset, rowCount = state.rowCount + 1)
+        val newSeenCount =
+          if (offset.timestamp == state.latestBacktracking.timestamp) 
state.latestBacktrackingSeenCount + 1 else 1
+
+        state.copy(
+          latestBacktracking = offset,
+          latestBacktrackingSeenCount = newSeenCount,
+          rowCount = state.rowCount + 1)
       } else {
         if (offset.timestamp.isBefore(state.latest.timestamp))
           throw new IllegalArgumentException(s"Unexpected offset [$offset] 
before latest [${state.latest}].")
 
+        if (log.isDebugEnabled()) {
+          if (state.latestBacktracking.seen.nonEmpty &&
+            
offset.timestamp.isAfter(state.latestBacktracking.timestamp.plus(firstBacktrackingQueryWindow)))
+            log.debug(
+              "{} next offset is outside the backtracking window, 
latestBacktracking: [{}], offset: [{}]",
+              logPrefix,
+              state.latestBacktracking,
+              offset)
+        }
+
         state.copy(latest = offset, rowCount = state.rowCount + 1)
       }
     }
@@ -351,7 +369,7 @@ import org.slf4j.Logger
     }
 
     def switchFromBacktracking(state: QueryState): Boolean = {
-      state.backtracking && state.rowCount < settings.bufferSize - 1
+      state.backtracking && state.rowCount < settings.bufferSize - 
state.backtrackingExpectFiltered
     }
 
     def nextQuery(state: QueryState): (QueryState, Option[Source[Envelope, 
NotUsed]]) = {
@@ -381,7 +399,8 @@ import org.slf4j.Logger
             queryCount = state.queryCount + 1,
             idleCount = newIdleCount,
             backtrackingCount = 1,
-            latestBacktracking = fromOffset)
+            latestBacktracking = fromOffset,
+            backtrackingExpectFiltered = state.latestBacktrackingSeenCount)
         } else if (switchFromBacktracking(state)) {
           // switch from backtracking
           state.copy(
@@ -398,7 +417,8 @@ import org.slf4j.Logger
             rowCountSinceBacktracking = state.rowCountSinceBacktracking + 
state.rowCount,
             queryCount = state.queryCount + 1,
             idleCount = newIdleCount,
-            backtrackingCount = newBacktrackingCount)
+            backtrackingCount = newBacktrackingCount,
+            backtrackingExpectFiltered = state.latestBacktrackingSeenCount)
         }
 
       val behindCurrentTime =
diff --git 
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/query/EventsBySliceBacktrackingSpec.scala
 
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/query/EventsBySliceBacktrackingSpec.scala
index 19cec16..3004862 100644
--- 
a/core/src/test/scala/org/apache/pekko/persistence/r2dbc/query/EventsBySliceBacktrackingSpec.scala
+++ 
b/core/src/test/scala/org/apache/pekko/persistence/r2dbc/query/EventsBySliceBacktrackingSpec.scala
@@ -14,6 +14,7 @@
 package org.apache.pekko.persistence.r2dbc.query
 
 import java.time.Instant
+import java.time.temporal.ChronoUnit
 
 import scala.concurrent.duration._
 import org.apache.pekko
@@ -23,6 +24,7 @@ import pekko.actor.typed.ActorSystem
 import pekko.persistence.query.NoOffset
 import pekko.persistence.query.Offset
 import pekko.persistence.query.PersistenceQuery
+import pekko.persistence.query.TimestampOffset
 import pekko.persistence.query.typed.EventEnvelope
 import pekko.persistence.r2dbc.Dialect
 import pekko.persistence.r2dbc.QuerySettings
@@ -44,9 +46,12 @@ import org.scalatest.wordspec.AnyWordSpecLike
 import org.slf4j.LoggerFactory
 
 object EventsBySliceBacktrackingSpec {
+  private val BufferSize = 10 // small buffer for testing
+
   private val config = ConfigFactory
-    .parseString("""
+    .parseString(s"""
     pekko.persistence.r2dbc.journal.publish-events = off
+    pekko.persistence.r2dbc.query.buffer-size = $BufferSize
     """)
     .withFallback(TestConfig.config)
 }
@@ -236,6 +241,52 @@ class EventsBySliceBacktrackingSpec
       // from normal query
       expect(result2.expectNext(), pid2, 3L, Some("e2-3"))
     }
+
+    "predict backtracking filtered events based on latest seen counts" in {
+      val entityType = nextEntityType()
+      val pid = nextPid(entityType)
+      val slice = query.sliceForPersistenceId(pid)
+      val sinkProbe = TestSink[EventEnvelope[String]]()(system.classicSystem)
+
+      // use times in the past well outside behind-current-time
+      val timeZero = 
InstantFactory.now().truncatedTo(ChronoUnit.SECONDS).minusSeconds(10 * 60)
+
+      // events around the buffer size (of 10) will share the same timestamp
+      // to test tracking of seen events that will be filtered on the next 
cycle
+
+      val numberOfCycles = 5 // loop through the grouped events this many times
+      val maxGroupSize = 7 // max size for groups of events that share the 
same timestamp
+      val eventsPerCycle = maxGroupSize * (maxGroupSize + 1) / 2 // each cycle 
has groups of size 1 to maxGroupSize
+      val numberOfEvents = numberOfCycles * eventsPerCycle // total number of 
events that will be generated
+
+      // generate a timestamp pattern like 1, 2, 2, 4, 4, 4, 7, 7, 7, 7, 11, 
11, 11, 11, 11, ...
+      for (cycle <- 1 to numberOfCycles) {
+        for (groupSize <- 1 to maxGroupSize) {
+          for (shift <- 1 to groupSize) {
+            val seqNr = (cycle - 1) * eventsPerCycle + (groupSize - 1) * 
groupSize / 2 + shift
+            val eventTime = seqNr - shift + 1
+            writeEvent(slice, pid, seqNr, timeZero.plusMillis(eventTime), 
s"event-$pid-$seqNr")
+          }
+        }
+      }
+
+      // start the query with a latest timestamp ahead of all events, to only 
run in backtracking mode
+      val latestOffset = TimestampOffset(timeZero.plusMillis(numberOfEvents + 
1), Map.empty)
+
+      val result: TestSubscriber.Probe[EventEnvelope[String]] =
+        query
+          .eventsBySlices[String](entityType, 0, persistenceExt.numberOfSlices 
- 1, latestOffset)
+          .runWith(sinkProbe)
+          .request(numberOfEvents)
+
+      // test that backtracking queries continued through all the written 
events, catching up with latest offset
+      for (seqNr <- 1 to numberOfEvents) {
+        val envelope = result.expectNext()
+        envelope.persistenceId shouldBe pid
+        envelope.sequenceNr shouldBe seqNr
+        envelope.source shouldBe EnvelopeOrigin.SourceBacktracking
+      }
+    }
   }
 
 }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to