Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 47 additions & 4 deletions internal/impl/gcp/enterprise/input_spanner_cdc.go
Original file line number Diff line number Diff line change
Expand Up @@ -183,9 +183,14 @@ type spannerCDCReader struct {

batching service.BatchPolicy
batcher *spannerPartitionBatcherFactory
res *service.Resources
Comment thread
squiidz marked this conversation as resolved.
Outdated
resCh chan asyncMessage
subscriber *changestreams.Subscriber
stopSig *shutdown.Signaller

// updateWatermark persists a partition watermark; set to the subscriber's
// UpdatePartitionWatermark in Connect, overridable in tests.
updateWatermark func(ctx context.Context, partitionToken string, ts time.Time) error
}

var _ service.BatchInput = (*spannerCDCReader)(nil)
Expand Down Expand Up @@ -217,24 +222,55 @@ func newSpannerCDCReader(conf spannerCDCInputConfig, batching service.BatchPolic
metrics: changestreams.NewMetrics(mgr.Metrics(), conf.StreamID),
batching: batching,
batcher: newSpannerPartitionBatcherFactory(batching, mgr),
res: mgr,
resCh: make(chan asyncMessage),
stopSig: shutdown.NewSignaller(),
}
}

// resetPartitionBatchers discards all cached partition batchers. Called on
// (re)connect: stale batchers hold rows and ack state from the previous
// subscriber session, and the re-read from the persisted watermarks
// re-delivers those rows anyway.
func (r *spannerCDCReader) resetPartitionBatchers() {
r.batcher = newSpannerPartitionBatcherFactory(r.batching, r.res)
Comment thread
squiidz marked this conversation as resolved.
Outdated
}

func (r *spannerCDCReader) emit(
ctx context.Context,
batcher *spannerPartitionBatcher,
partitionToken string,
msg service.MessageBatch,
commitTimestamp time.Time,
) (*ack.Once, error) {
if len(msg) == 0 {
return nil, nil
}
// A zero commitTimestamp means "mid-record, do not advance": substitute
// the last known-safe watermark so an out-of-order resolve can never
// regress or stall behind it. Per-partition callbacks are serialized, so
// this needs no locking.
if commitTimestamp.IsZero() {
commitTimestamp = batcher.lastWatermark
} else {
batcher.lastWatermark = commitTimestamp
}
resolveFn, err := batcher.cp.Track(ctx, commitTimestamp, int64(len(msg)))
if err != nil {
return nil, fmt.Errorf("tracking watermark checkpoint: %w", err)
}
ackOnce := ack.NewOnce(func(ctx context.Context) error {
// Only the resolved (contiguous-prefix) watermark is safe to persist:
// a batch acked out of order must not advance the watermark past
// still-unacked earlier batches, or a crash in that window would skip
// their records on restart.
resolved := resolveFn()
if resolved == nil || resolved.IsZero() {
return nil
}
// If we processed the message and failed to update the watermark, we
// would try to update it on the next message, no need to return an error here.
if err := r.subscriber.UpdatePartitionWatermark(ctx, partitionToken, commitTimestamp); err != nil {
if err := r.updateWatermark(ctx, partitionToken, *resolved); err != nil {
r.log.Errorf("%s: failed to update watermark: %v", partitionToken, err)
}
return nil
Expand Down Expand Up @@ -268,7 +304,7 @@ func (r *spannerCDCReader) onDataChangeRecord(ctx context.Context, partitionToke
if err != nil {
return err
}
ack, err := r.emit(ctx, partitionToken, msg, ts)
ack, err := r.emit(ctx, batcher, partitionToken, msg, ts)
if err != nil {
return err
}
Expand All @@ -289,7 +325,7 @@ func (r *spannerCDCReader) onDataChangeRecord(ctx context.Context, partitionToke
if err != nil {
return err
}
ack, err := r.emit(ctx, partitionToken, msg, ts)
ack, err := r.emit(ctx, batcher, partitionToken, msg, ts)
if err != nil {
return err
}
Expand All @@ -300,7 +336,7 @@ func (r *spannerCDCReader) onDataChangeRecord(ctx context.Context, partitionToke

iter := batcher.MaybeFlushWith(dcr)
for mb, ts := range iter.Iter(ctx) {
ack, err := r.emit(ctx, partitionToken, mb, ts)
ack, err := r.emit(ctx, batcher, partitionToken, mb, ts)
if err != nil {
return err
}
Expand All @@ -327,11 +363,18 @@ func (r *spannerCDCReader) Connect(ctx context.Context) error {
cb = p.onDataChangeRecord
}

// Discard any partition batchers from a previous subscriber session:
// their buffered rows were never acked and will be re-read from the
// persisted watermarks; reusing them would duplicate rows into mixed
// batches and misalign the ack tracker.
r.resetPartitionBatchers()
Comment thread
squiidz marked this conversation as resolved.
Outdated

var err error
r.subscriber, err = changestreams.NewSubscriber(ctx, r.conf.Config, cb, r.log, r.metrics)
if err != nil {
return fmt.Errorf("create Spanner change stream reader: %w", err)
}
r.updateWatermark = r.subscriber.UpdatePartitionWatermark

if err := r.subscriber.Setup(ctx); err != nil {
return fmt.Errorf("setup Spanner change stream reader: %w", err)
Expand Down
172 changes: 172 additions & 0 deletions internal/impl/gcp/enterprise/input_spanner_cdc_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
// Copyright 2026 Redpanda Data, Inc.
//
// Licensed as a Redpanda Enterprise file under the Redpanda Community
// License (the "License"); you may not use this file except in compliance with
// the License. You may obtain a copy of the License at
//
// https://github.com/redpanda-data/connect/blob/main/licenses/rcl.md

package enterprise

import (
"context"
"sync"
"testing"
"time"

"github.com/stretchr/testify/require"

"github.com/redpanda-data/benthos/v4/public/service"
)

// watermarkRecorder is a test seam for spannerCDCReader.updateWatermark.
type watermarkRecorder struct {
mu sync.Mutex
writes []time.Time
}

func (w *watermarkRecorder) update(_ context.Context, _ string, ts time.Time) error {
w.mu.Lock()
defer w.mu.Unlock()
w.writes = append(w.writes, ts)
return nil
}

func (w *watermarkRecorder) recorded() []time.Time {
w.mu.Lock()
defer w.mu.Unlock()
return append([]time.Time(nil), w.writes...)
}

func newTestSpannerReader(t *testing.T) (*spannerCDCReader, *spannerPartitionBatcher, *watermarkRecorder) {
t.Helper()

r := newSpannerCDCReader(spannerCDCInputConfig{}, service.BatchPolicy{Count: 1}, service.MockResources())
rec := &watermarkRecorder{}
r.updateWatermark = rec.update

// Drain emitted messages so emit's channel send never blocks; acks are
// driven explicitly by the tests via the returned ack.Once.
go func() {
for {
select {
case <-t.Context().Done():
return
case <-r.resCh:
}
}
}()

batcher, _, err := r.batcher.forPartition("p1")
require.NoError(t, err)
return r, batcher, rec
}

func testBatch() service.MessageBatch {
return service.MessageBatch{service.NewMessage([]byte("{}"))}
}

func TestSpannerEmitOrderedWatermarks(t *testing.T) {
t1 := time.Date(2026, 8, 6, 10, 0, 0, 0, time.UTC)
t2 := t1.Add(time.Second)
t3 := t2.Add(time.Second)

t.Run("out-of-order acks never advance past unacked batches", func(t *testing.T) {
ctx := t.Context()
r, batcher, rec := newTestSpannerReader(t)

ack1, err := r.emit(ctx, batcher, "p1", testBatch(), t1)
require.NoError(t, err)
ack2, err := r.emit(ctx, batcher, "p1", testBatch(), t2)
require.NoError(t, err)

// Acking only the LATER batch must not write anything: its records
// are durable but the earlier batch's are not. (The pre-fix code
// wrote t2 here - the data-loss window.)
require.NoError(t, ack2.Ack(ctx, nil))
require.Empty(t, rec.recorded(), "watermark must not advance past a still-unacked earlier batch")

// Acking the earlier batch resolves the full prefix: one write, t2.
require.NoError(t, ack1.Ack(ctx, nil))
require.Equal(t, []time.Time{t2}, rec.recorded())
})

t.Run("in-order acks advance incrementally", func(t *testing.T) {
ctx := t.Context()
r, batcher, rec := newTestSpannerReader(t)

ack1, err := r.emit(ctx, batcher, "p1", testBatch(), t1)
require.NoError(t, err)
ack2, err := r.emit(ctx, batcher, "p1", testBatch(), t2)
require.NoError(t, err)

require.NoError(t, ack1.Ack(ctx, nil))
require.NoError(t, ack2.Ack(ctx, nil))
require.Equal(t, []time.Time{t1, t2}, rec.recorded())
})

t.Run("zero-watermark batches carry the last safe watermark forward", func(t *testing.T) {
ctx := t.Context()
r, batcher, rec := newTestSpannerReader(t)

ack1, err := r.emit(ctx, batcher, "p1", testBatch(), t1)
require.NoError(t, err)
require.NoError(t, ack1.Ack(ctx, nil))
require.Equal(t, []time.Time{t1}, rec.recorded())

// b2 is a mid-record flush (zero watermark), b3 completes a record.
ack2, err := r.emit(ctx, batcher, "p1", testBatch(), time.Time{})
require.NoError(t, err)
ack3, err := r.emit(ctx, batcher, "p1", testBatch(), t3)
require.NoError(t, err)

// Acking b3 alone must not advance past t1 (b2 is unacked; a repeat
// write of the current safe watermark t1 is fine and idempotent).
// The pre-fix code wrote t3 here - the data-loss window.
require.NoError(t, ack3.Ack(ctx, nil))
for _, w := range rec.recorded() {
require.False(t, w.After(t1), "watermark advanced past an unacked batch: %v", w)
}
// Acking b2 resolves the full prefix through b3.
require.NoError(t, ack2.Ack(ctx, nil))
writes := rec.recorded()
require.Equal(t, t3, writes[len(writes)-1])
})

t.Run("zero-watermark batch acked alone re-writes only the safe watermark", func(t *testing.T) {
ctx := t.Context()
r, batcher, rec := newTestSpannerReader(t)

ack1, err := r.emit(ctx, batcher, "p1", testBatch(), t1)
require.NoError(t, err)
ack2, err := r.emit(ctx, batcher, "p1", testBatch(), time.Time{})
require.NoError(t, err)

require.NoError(t, ack1.Ack(ctx, nil))
require.NoError(t, ack2.Ack(ctx, nil))

// The carried-forward value is t1; writes never regress and never
// mention a timestamp past the last completed record.
writes := rec.recorded()
require.NotEmpty(t, writes)
for _, w := range writes {
require.Equal(t, t1, w)
}
})
}

func TestSpannerConnectResetsPartitionBatchers(t *testing.T) {
r, batcher, _ := newTestSpannerReader(t)
_ = batcher

// Same token returns the cached batcher before reset...
_, existed, err := r.batcher.forPartition("p1")
require.NoError(t, err)
require.True(t, existed)

// ...and a fresh one after the reset performed on (re)connect.
r.resetPartitionBatchers()
_, existed, err = r.batcher.forPartition("p1")
require.NoError(t, err)
require.False(t, existed, "reconnect must not reuse stale partition batcher state")
}
16 changes: 16 additions & 0 deletions internal/impl/gcp/enterprise/input_spanner_partition_batcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ import (
"sync"
"time"

"github.com/Jeffail/checkpoint"

"github.com/redpanda-data/benthos/v4/public/service"
"github.com/redpanda-data/connect/v4/internal/ack"
"github.com/redpanda-data/connect/v4/internal/impl/gcp/enterprise/changestreams"
Expand Down Expand Up @@ -103,12 +105,25 @@ func (s *spannerPartitionBatchIter) Err() error {
return s.err
}

// spannerCheckpointCap bounds the number of in-flight (unacked) batches
// tracked per partition; Track blocks once reached, applying backpressure.
const spannerCheckpointCap = 1024
Comment thread
squiidz marked this conversation as resolved.
Outdated

type spannerPartitionBatcher struct {
batcher *service.Batcher
last *changestreams.DataChangeRecord
period *time.Timer
acks []*ack.Once
rm func()

// cp orders in-flight batches so the partition watermark only ever
// advances to a commit timestamp once every batch at or below it has been
// acked (see spannerCDCReader.emit).
cp *checkpoint.Capped[time.Time]
// lastWatermark is the most recent non-zero watermark emitted for this
// partition; mid-record (zero watermark) batches carry it forward. Only
// accessed from the partition's serialized callback.
lastWatermark time.Time
}

func (s *spannerPartitionBatcher) MaybeFlushWith(dcr *changestreams.DataChangeRecord) *spannerPartitionBatchIter {
Expand Down Expand Up @@ -202,6 +217,7 @@ func (f *spannerPartitionBatcherFactory) forPartition(partitionToken string) (*s

spb = &spannerPartitionBatcher{
batcher: b,
cp: checkpoint.NewCapped[time.Time](spannerCheckpointCap),
rm: func() {
f.mu.Lock()
delete(f.partitions, partitionToken)
Expand Down