Skip to content

aws_dynamodb_cdc: add snapshot_mode incremental (ENG-1369) - #4886

Open
squiidz wants to merge 1 commit into
mainfrom
eng-1369-dynamodb-incremental-snapshot
Open

squiidz wants to merge 1 commit into
mainfrom
eng-1369-dynamodb-incremental-snapshot

Conversation

@squiidz

@squiidz squiidz commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

What

Adds snapshot_mode: incremental to aws_dynamodb_cdc. The input backfills each table page by page while it streams. This works in single-table and multi-table mode; in multi-table mode it also covers tables found by periodic discovery. The guarantee is that no snapshot item leaves the input after a newer stream event for the same key, so the last value per key that the input emits matches the table.

Jira: ENG-1369

Why

snapshot_and_cdc already streams while it scans, but it has two gaps:

  • A stale scan item can overtake a newer change. The Scan reads V1, the item changes to V2, and the stream delivers V2 first. V1 is still emitted afterwards, so downstream ends on V1. The dedupe buffer also switches itself off once it overflows.
  • No snapshot in multi-table mode. Existing rows of discovered tables are never read.

snapshot_and_cdc is unchanged. incremental is a new mode.

How

Each Scan page, per segment, is a window:

  1. The page opens before its Scan request. The Scan uses ConsistentRead.
  2. Shard readers check every record's key against open pages before enqueueing the record. A match drops the held item. The stream will deliver that key's current value, because a new image is required.
  3. Releasing a page sends it to the message channel while holding the same lock. For any key, either the stream record drops the item, or the item is enqueued first.

This guarantee does not depend on clocks. On top of it, a per-shard progress tracker holds each page until every shard has streamed past the read. That keeps snapshot rows behind older stream events. If the tracker misjudges, the only effect is a brief reorder; the final value is still correct.

Released items go through the existing ack-gated snapshot path, so segment progress is checkpointed only after acks. In multi-table mode, tables are backfilled one at a time. Staleness is checked at connect and at discovery, before the table's readers start.

Config

  • snapshot_mode: incremental, allowed with every table_discovery_mode.
  • snapshot_watermark_margin (advanced, default 2s): extra hold time for clock skew and stream publication delay. It only affects ordering.
  • Requires a NEW_IMAGE or NEW_AND_OLD_IMAGES stream view. Single-table mode fails Connect; multi-table mode logs the table and skips its backfill, but keeps streaming it.
  • Strongly consistent Scans cost twice the RCU.
  • A table that was only ever streamed is backfilled once, the first time incremental is enabled.

New metrics:

  • dynamodb_cdc_snapshot_window_held_items
  • dynamodb_cdc_snapshot_window_dropped_items
  • dynamodb_cdc_snapshot_window_wait
  • dynamodb_cdc_snapshot_backfill_tables_pending

Signals (snapshot on demand)

A new optional signal_table_name names a DynamoDB table with a NEW_IMAGE or NEW_AND_OLD_IMAGES stream that carries control signals. The signal format and decoding match postgres_cdc, using the shared internal/replication helpers from #4919.

  • Signal items. A signal is an INSERT of an item with id (a new value for every signal), type and data (a JSON string). Re-sending an existing id becomes a MODIFY, which is ignored with a Warn.
  • log writes data.message to the connector's log.
  • snapshot-execute with {"tables": [...]} backfills each named table once. A table that is already complete is skipped with a Warn. If any named table is unwatched or uses the wrong stream view, the whole request is rejected.
  • Signal-only mode. With snapshot_mode: incremental and signal_table_name set, backfills happen only on request. Connect resumes only:
    • tables that were in the middle of a backfill;
    • tables with a pending request;
    • completed tables whose stream checkpoint went stale. These are reset and re-backfilled automatically.
  • Durability. A snapshot#requested checkpoint row is written before the signal record can be acked. A failed write is retried, and the signal is held back until it succeeds. A stale reset records the request first, so a crash at any point does not lose it.
  • Downstream. Signal records are emitted downstream like any other record. The docs show a mapping that filters them out.
  • Supported modes. signal_table_name works with snapshot_mode none or incremental. A signal table makes the input multi-table, so snapshot_only and snapshot_and_cdc are rejected.

Integration test TestIntegrationDynamoDBCDCSnapshotSignal, run against DynamoDB Local:

  • Before the signal, no snapshot checkpoint rows exist.
  • After the signal, the table is backfilled while a concurrent writer updates it, with 0 safety violations and final values matching the table.
  • A repeat signal writes no request and emits no further snapshot items.

Behaviour changes outside the new mode

  • SnapshotProgress now reads the checkpoint table with strongly consistent reads in all modes. These reads are rare.
  • The snapshot_and_cdc dedupe buffer is no longer allocated for snapshot_only, which never streams.

Known limitations

  • Throughput. Each segment holds one page for a few seconds, so defaults give roughly 20 items/s. The docs recommend snapshot_segments: 10 and snapshot_batch_size: 1000 for large tables. Holding several pages per segment is a possible follow-up.
  • Retry after failure. A failed backfill is logged, the snapshot state metric shows failed, and streaming continues. It is retried on the next restart.
  • Connect time. In multi-table mode, Connect runs a staleness check per table, so it is slower with many tables.
  • DynamoDB Local. It rounds ApproximateCreationDateTime down to the minute, so pages are released late locally, mostly after writes go quiet.

Testing

  • Unit tests cover the window (including a concurrency test on the release/touch ordering), the shard tracker, the key encoding, the scanner hook, the backfiller (completion only after acks, cancel then rerun, touch after hold, keyless item), the reset ordering under stale reads, the multi-table queue, and the connect paths.
  • go test -race passes for the package. task lint is clean.
  • Integration test (TestIntegrationDynamoDBCDCIncrementalSnapshot): 500 items, with a concurrent writer doing ADD v 1 for 10s. It checks per key that a READ value is never lower than a value already emitted for that key, and that the final values match the table.
    • Against DynamoDB Local: passes, with 400 READ events and 0 violations.
    • Against real AWS (DYNAMODB_CDC_REAL_AWS=1), which also asserts per-key ordering: not yet run.

@squiidz
squiidz marked this pull request as draft September 29, 2026 16:05
Comment thread internal/impl/aws/dynamodb/input_cdc_integration_test.go Outdated
Comment thread internal/impl/aws/dynamodb/input_cdc_integration_test.go
@squiidz
squiidz force-pushed the eng-1369-dynamodb-incremental-snapshot branch from a05889a to d24b2f9 Compare October 1, 2026 12:43
Comment thread internal/impl/aws/dynamodb/snapshot_watermark.go Outdated
Comment thread internal/impl/aws/dynamodb/snapshot_incremental_backfill_test.go Outdated
@squiidz
squiidz force-pushed the eng-1369-dynamodb-incremental-snapshot branch from d24b2f9 to c446da9 Compare October 1, 2026 13:36
@squiidz
squiidz marked this pull request as ready for review October 2, 2026 17:32
Comment on lines +76 to +86
func snapshotPKs(t *testing.T, batch service.MessageBatch) []string {
t.Helper()
var pks []string
for _, msg := range batch {
s, err := msg.AsStructured()
require.NoError(t, err)
img := s.(map[string]any)["dynamodb"].(map[string]any)["newImage"].(map[string]any)
pks = append(pks, img["pk"].(string))
}
return pks
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

require is called from goroutines other than the test goroutine (tester pattern: "Do not use require inside assert.Eventually. require calls FailNow() which panics when called from a non-test goroutine")

snapshotPKs calls require.NoError. Most callers run it inside the background goroutines that drain d.msgChan. That happens in TestIncrementalBackfillEmitsUntouchedItemsAndCompletes, TestIncrementalBackfillHoldsUntilCaughtUp, TestIncrementalBackfillTouchAfterHoldDrops, TestIncrementalBackfillCancelResetsWindowForRerun, TestIncrementalBackfillResetThenStaleReadsStillBackfills, and in TestIncrementalBackfillResumesFromAckedPosition in snapshot_incremental_single_test.go.

When FailNow is called off the test goroutine, it only exits the drain goroutine. After that, nothing reads msgChan. runIncrementalBackfill then blocks in handleSnapshotBatch, so the test hangs until a timeout or the package deadline. It does not fail cleanly with the real error.

Fix: give the helper a non-fatal form for these drain goroutines, for example one that uses assert.NoError and returns, or one that returns an error. Keep require only in the calls made from the test goroutine.

See the Polling section of the tester guidance and the CLAUDE.md skills table.

@squiidz
squiidz force-pushed the eng-1369-dynamodb-incremental-snapshot branch from c446da9 to 70c3130 Compare October 5, 2026 17:18
Comment thread internal/impl/aws/dynamodb/input_cdc.go Outdated
dciFieldSnapshotBufferSize = "snapshot_buffer_size"
dciFieldSnapshotWatermarkMargin = "snapshot_watermark_margin"
dciFieldSnapshotIdleShardGrace = "snapshot_idle_shard_grace"
dciFieldSignalTable = "signal_table"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CDC field naming doesn't match the fleet (CONTRIBUTING §5, §5.3.3)

The commit message says the signal design is "matching postgres_cdc", but postgres_cdc names this field signal_table_name (input_pg_stream.go#L54-L56). This PR adds signal_table, so the same feature gets a different name here.

§5 says to mirror "the shape the fleet already establishes", and §5.3.3 says "Do not invent bespoke names where a canonical one already exists." Since this field is new, renaming it now is free; changing it after release would need a deprecation (§5.3.2). Use signal_table_name so the two connectors match. Update the constant, the docs, and the validation and error messages to the new name.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Renamed to signal_table_name in 347d325, matching postgres_cdc. Constant, docs, validation and error messages updated.

Comment on lines +1987 to +1992
require.Eventually(t, func() bool {
mu.Lock()
seen, sigs := len(maxSeen), signalRecords
mu.Unlock()
return seen == itemCount && sigs >= 1 && len(snapshotRows("snapshot#complete")) == 1
}, 5*time.Minute, 500*time.Millisecond, "timed out waiting for the backfill and the signal record")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

require called inside an Eventually condition (tester pattern: Polling)

This condition calls snapshotRows, which runs require.NoError(t, err) (L1921-L1937). testify runs the Eventually condition on a separate goroutine. If the checkpoint Query fails there, FailNow runs off the test goroutine, so the test aborts incorrectly instead of retrying the poll.

The project test patterns say: "Do not use require inside assert.Eventually. require calls FailNow() which panics when called from a non-test goroutine. Use assert or return bool."

Suggested fix: give the polling path a variant of snapshotRows that returns ([]string, error), and have the condition return false on error.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 347d325: the polling condition now uses snapshotRowsErr, which returns ([]string, error), and the condition returns false on error. snapshotRows keeps require and is only called on the test goroutine.

Comment on lines +138 to +143
func (d *dynamoDBCDCInput) retrySnapshotSignal(ctx context.Context, log *service.Logger, data []byte) bool {
boff := backoff.NewExponentialBackOff()
boff.InitialInterval = 200 * time.Millisecond
boff.MaxInterval = 5 * time.Second
boff.MaxElapsedTime = 0 // Never give up
for {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hardcoded retry durations (godev: Configurable Time Parameters)

This retry backoff is hardcoded: InitialInterval = 200ms and MaxInterval = 5s. The project Go patterns say: "Every time-related value (timeouts, backoffs, intervals, retry delays) must be exposed as a YAML-configurable field. Do not hardcode durations."

The same rule applies to two other new hardcoded durations in this PR:

  • the fixed time.Second added in canRelease
  • releaseTick = 250ms

Both are in snapshot_incremental.go#L113-L123.

Suggested fix: expose these as advanced fields, or reuse an existing config duration where one fits. For example, throttle_backoff could serve as the signal retry backoff, and snapshot_watermark_margin could absorb the fixed 1s.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Kept as documented constants in 347d325. The time.Second in canRelease is a protocol constant rather than a tuning knob: Streams round ApproximateCreationDateTime down to the second, so a shard seen up to second S may still deliver records from later in S. The user-tunable slack is snapshot_watermark_margin. The signal retry backoff mirrors the reader's existing hardcoded GetRecords backoff and only applies while a checkpoint write for a signal is failing. releaseTick was settled in the earlier round: it is a fallback poll behind progress-driven releases.

Backfill tables page by page while streaming, in single- and multi-table
mode, including tables found by discovery in multi-table mode. No snapshot
item leaves the input after a newer stream event for the same key.

Each Scan page opens a window before its request (with ConsistentRead).
A stream record for a held key drops the item, and release shares a lock
with that check, so the guarantee does not depend on clocks. A per-shard
progress tracker additionally holds pages until every shard has streamed
past the read, keeping snapshot rows behind older stream events.

Also:
- the window uses an injective key encoding (the dedupe key builders
  collapse binary keys and do not escape separators)
- stale-checkpoint recovery resets snapshot progress segments first and
  the marker last, and SnapshotProgress reads are strongly consistent
- multi-table backfills are prepared at connect and discovery, before
  the table's coordinator starts, and run one table at a time
- a failed backfill is logged and streaming continues

Signals: a new `signal_table_name` (a DynamoDB table with a stream) carries
control signals, using the shared internal/replication decoding. A
`snapshot-execute` INSERT, `{"tables": [...]}`, backfills each named
table once, matching postgres_cdc. With `signal_table_name` set, incremental
backfills are signal-only: connect resumes in-progress, requested or
stale tables, and waits for a signal otherwise. Requests are recorded
durably (`snapshot#requested`) before the signal record can be acked.
@squiidz
squiidz force-pushed the eng-1369-dynamodb-incremental-snapshot branch from 70c3130 to 347d325 Compare October 5, 2026 18:05
Comment on lines +143 to +146
boff := backoff.NewExponentialBackOff()
boff.InitialInterval = 200 * time.Millisecond
boff.MaxInterval = 5 * time.Second
boff.MaxElapsedTime = 0 // Never give up

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hardcoded retry durations break the godev "Configurable Time Parameters" rule

The new signal retry backoff hardcodes its durations (InitialInterval = 200ms, MaxInterval = 5s). The project Go patterns say: "Every time-related value (timeouts, backoffs, intervals, retry delays) must be exposed as a YAML-configurable field. Do not hardcode durations."

The comment says this copies the reader's hardcoded GetRecords backoff. That older code doesn't exempt new code from the rule. Two ways to fix it:

  • Use the existing throttle_backoff field.
  • Add an advanced duration field for checkpoint write retries.

The new releaseTick fallback poll interval (snapshot_incremental.go, var releaseTick = 250 * time.Millisecond) is another hardcoded interval of the same kind.

Ref: CLAUDE.md → godev agent, "Configurable Time Parameters".

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Keeping both as documented constants; same reasoning as the earlier thread on these durations. The signal retry backoff only runs while a checkpoint read or write for a signal is failing, never gives up, and holds the signal until the write succeeds, so the interval affects only how quickly a recovered checkpoint table is noticed, not correctness or throughput. releaseTick is a fallback poll: releases are triggered by stream progress, so changing it alters release latency by at most one tick. Neither is a knob an operator would usefully tune, and adding fields for them would widen the config surface without a use case. User-facing timing is configurable through snapshot_watermark_margin, snapshot_idle_shard_grace and snapshot_throttle.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant