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]