Skip to content

Change streams

fdyno implements the DynamoDB Streams read API over change records stored in FoundationDB. A record is written in the same FoundationDB transaction as its item and secondary-index changes. This gives atomic change capture, but consumers still own checkpoint durability, duplicate processing, retention lag, and sink idempotency.

Atomic capture and sequence order

For a stream-enabled table, a state-changing write adds a CDC record under the table's v subspace:

flowchart LR
    W["Item mutation"] --> T
    subgraph T["one FoundationDB transaction"]
      I["base item and chunks"]
      X["secondary-index entries"]
      V["CDC record at a versionstamped key"]
    end

The key is written with FoundationDB SetVersionstampedKey. Its 12-byte versionstamp contains a 10-byte commit version and a two-byte user version. fdyno encodes those bytes as a fixed-width decimal DynamoDB SequenceNumber.

The resulting storage guarantees are:

  • A committed state change and its record are visible together. A transaction cannot commit the item while omitting its record, or vice versa.
  • Sequence numbers are unique and increase in FoundationDB commit order within a shard.
  • All mutations for the same partition-key value route to one shard, so their sequence order is their commit order.
  • A multi-item FoundationDB transaction gives its records one commit version and distinct user versions.

A put that writes bytes equal to the current item and a delete of an absent item do not create a record. A retry of a non-idempotent update can commit a second state change, in which case the stream correctly contains two records. Write-side atomicity is not an exactly-once consumer-delivery guarantee; see Reliable writes.

Do not order records by ApproximateCreationDateTime. That field is generated from the fdyno process wall clock before commit and can skew across replicas. Use sequence numbers for order.

Choose the shard count when enabling a stream

DYNODB_STREAM_SHARDS supplies the default shard count for a newly enabled stream generation. Valid values are 1 through 100. Invalid values use one shard. fdyno stores the selected count, topology version, and random generation ID in table metadata.

Every writer, reader, iterator, description call, chunk lookup, and retention worker uses that persisted topology. Replicas can have different environment defaults without rerouting an active generation. Use the same default on all replicas so a create or enable request has the same result regardless of which process receives it.

With one shard, a table has one commit-ordered log. With more than one, fdyno hashes canonical binary partition-key bytes modulo the persisted count. Consumers can read shards in parallel. Records for different partition keys have no single delivery order.

Changing the count starts an empty generation

An active generation has an immutable shard count. To select another count, disable the stream, set the new default, and re-enable it. Re-enabling clears the disabled generation's records, creates new shard IDs, and invalidates its iterators.

Static shards do not split or merge as load changes and have no parent/child lineage. DescribeStream returns the configured fixed set. Its reported starting and ending sequence-range metadata is simplified and is not a measured oldest/newest record watermark. ShardFilter=CHILD_SHARDS therefore returns no child shards.

Increasing the count only adds consumer scan parallelism. It does not increase FoundationDB commit capacity, and it cannot distribute a single hot partition key without violating per-key order.

Stream API and iterator positions

fdyno exposes ListStreams, DescribeStream, GetShardIterator, and GetRecords. The supported iterator types are:

Iterator type Initial position
TRIM_HORIZON First record currently retained in the shard.
LATEST Immediately after the newest record visible when the iterator is created.
AT_SEQUENCE_NUMBER The key represented by that sequence number.
AFTER_SEQUENCE_NUMBER Immediately after that key.

An iterator is an opaque base64 cursor containing stream ARN, shard ID, and the next FoundationDB key. It is not server-side state, has no signature, and does not expire on a timer. GetRecords reads at most 1,000 records, or the lower requested limit, and returns a cursor after the last scanned key. Polling an empty enabled shard returns another iterator at the same position.

When a stream is disabled, no new records are written. A consumer can drain retained records; once a disabled shard returns fewer than the requested page limit, NextShardIterator is empty.

Disabling retains that generation long enough for consumers to drain stored records. Re-enabling atomically clears its v and vc ranges before publishing a new generation identity and shard set. The disabled ARN and its iterators stop resolving. fdyno does not provide DynamoDB's overlapping access to the disabled generation.

View types, storage format, and write amplification

All four view types are accepted:

View type Stored payload
KEYS_ONLY Primary key only.
NEW_IMAGE Post-mutation image.
OLD_IMAGE Pre-mutation image.
NEW_AND_OLD_IMAGES Both images.

DynamoDB Streams responses remain JSON. FoundationDB values use the deterministic, versioned FDYNOC01 binary record. Keys and images are length-framed values from the existing item codec.

A record through 10,000 bytes is stored directly at its ordered v key. A larger record stores the existing versioned manifest there and raw binary chunks under vc. Readers validate chunk count, exact size, SHA-256, and the record and item frames. Stored JSON-formatted direct and chunked records are also readable and trimmable.

Chunking removes the single-value boundary but not FoundationDB's transaction limit. NEW_AND_OLD_IMAGES has the largest write amplification. Item, index, record manifest, and all record chunks still commit in one transaction.

Choose the smallest view that supplies the consumer's invariant. KEYS_ONLY lowers storage and transaction amplification but forces consumers to read current table state, which may already be newer than the event. Image-bearing views preserve the change context at a higher storage and write cost.

TTL expiry writes a REMOVE record when the item is transactionally deleted, with a service identity. Read paths can hide an expired item before background deletion, so “not returned by reads” and “TTL remove record emitted” need not happen at the same time.

Consumer checkpoints and delivery semantics

fdyno does not maintain consumer groups, acknowledgements, or checkpoints. The consumer must keep one durable checkpoint per stream ARN and shard.

Use this processing pattern:

  1. Obtain TRIM_HORIZON, LATEST, or an AFTER_SEQUENCE_NUMBER iterator from the durable checkpoint.
  2. Read a bounded page.
  3. Apply each record to an idempotent sink, using stream ARN, shard ID, and sequence number as a deduplication identity where appropriate.
  4. Persist the last successfully applied sequence number only after the sink effect is durable.
  5. Resume with AFTER_SEQUENCE_NUMBER after restart.

Checkpoint-before-effect can lose work. Effect-before-checkpoint can repeat work after a crash. Unless the sink and checkpoint share a transaction, choose at-least-once processing and make the sink idempotent.

Sequence numbers are stable when older keys are removed; trimming does not renumber surviving records. However, fdyno does not return an expired-iterator or trimmed-data error when a cursor points before the oldest retained key. The range scan naturally starts at the first surviving key. A lagging consumer can therefore skip trimmed records without an explicit gap signal. Monitor lag externally and retain enough history for the worst expected outage.

Retention and trimming

Without DYNODB_CDC_TRIM_INTERVAL, change records remain indefinitely. The default configuration therefore favors not deleting consumer history, but storage grows without bound.

When enabled, the trimmer uses:

  • DYNODB_CDC_TRIM_INTERVAL for scheduling;
  • DYNODB_CDC_RETENTION, defaulting to 24 hours, for the cutoff; and
  • DYNODB_CDC_TRIM_MAX_SCAN, defaulting to 1,000, per table and shard each pass.

The worker scans from the oldest key and deletes only a contiguous prefix whose stored ApproximateCreationDateTime is strictly older than the cutoff. It stops at an undecodable, undated, or in-window record rather than deleting uncertain data. Transactions run at FoundationDB batch priority and replicas coordinate with a cooperative lease.

Retention is based on fdyno host wall clocks, not the FoundationDB commit version. Clock skew can retain records too long or trim them early. Zero or negative retention can remove current history. Synchronize clocks and validate duration values before startup.

A bounded trimmer can fall behind. For each shard, its maximum examined records per interval must exceed the sustained arrival rate plus catch-up demand; otherwise the oldest age continues to increase. More shards multiply each table's per-pass scan allowance but also add worker transactions.

Capacity and observability

Stream capacity must include:

  • at least one FoundationDB record value per state-changing mutation, plus chunks when the encoded record exceeds the direct-value bound;
  • key and binary record-encoding overhead;
  • one or two item images depending on view type;
  • retained record count during normal lag and consumer outage; and
  • read load from every polling consumer.

GET /metrics and GET /metrics/prometheus expose process-local cdc-trimmer lease/tick and deleted-item counters plus surfaced error attempts. Zero reported errors does not prove a successful trim pass. Neither endpoint reports records written, bytes retained, oldest record age, consumer lag, per-shard read rate, or all corruption events.

GetRecords skips an undecodable stored key or value and advances its cursor. It does not increment an API error counter. Detect storage corruption with FoundationDB monitoring and application sequence checks.

Recommended external signals are oldest retained timestamp per shard, latest seen sequence, durable consumer checkpoint, records/bytes processed, duplicate count, sink failures, and wall-clock lag. Keep table/shard labels controlled to avoid an unbounded metrics cardinality problem.

Incident checks

When a consumer stops receiving records:

  1. Confirm the table's stream is enabled and the consumer uses LatestStreamArn.
  2. Use DescribeStream to read the persisted shards for the active generation. Do not infer the active count from a process environment variable.
  3. Describe all shards and verify the consumer owns and polls each fixed shard.
  4. Check whether the iterator was created with LATEST after the expected writes. Such an iterator starts after them.
  5. Read from TRIM_HORIZON in a diagnostic consumer without advancing the production checkpoint.
  6. Check FoundationDB readiness and transaction errors on the write path; atomic capture means a successfully committed state change should have its record.
  7. Determine whether retention passed the consumer checkpoint. Absence of an API trim error does not prove continuity.
  8. Check trimmer counters, clock synchronization, retention configuration, and whether an undecodable oldest record has blocked prefix trimming.

When stream-enabled writes fail but the same writes succeed with streams disabled, check record value size and total transaction amplification, especially with NEW_AND_OLD_IMAGES and multiple secondary indexes.