DynamoDB Streams Shard Architecture
Shards enforce ordering, isolation, and consumer assignment in DynamoDB Streams.

DynamoDB Streams is a change data capture mechanism. It records every item-level modification to a DynamoDB table, whether an INSERT, a MODIFY, or a REMOVE, and arranges those records into a time-ordered sequence organized by shard. The shard is the structural unit that governs ordering, isolation, and which consumer gets assigned to which slice of change data. AWS's own documentation on the subject establishes two guarantees that only make sense once the shard is understood as the base unit: each stream record appears exactly once in the stream, and records for a given item arrive in the same sequence as the modifications that produced them. Engineers who approach DynamoDB Streams as though it were a generic message queue, without grasping what a shard physically represents, end up writing consumers that lose ordering guarantees, miss records after a shard splits, or run straight into the 24-hour retention wall with no warning. The shard architecture is the structural foundation for every reliability property the stream offers, and nothing downstream (view types, retention, consumer semantics) can be reasoned about correctly without it.
The 1:1 partition-to-shard mapping
The central fact of DynamoDB Streams architecture is a strict 1:1 relationship between a table partition and its open stream shard. Every write to a given partition flows into that partition's dedicated open shard, and no other partition writes into that same shard. This holds for one open, writable shard per partition at a time, not for every shard that has ever existed in the stream's history.
That mapping enforces three things at once. Ordering within a partition is guaranteed: records for a given partition key arrive in the shard in the exact sequence in which they were written to the table. Ordering across partitions is not guaranteed in the same way. A consumer that needs a single global order across every partition in the table cannot get it from DynamoDB Streams alone, because each partition's changes live in their own isolated shard. Isolation between partitions is structural: the two partitions were never writing to the same shard in the first place, so no mechanism exists by which their records could interleave inside a single shard.
The mapping extends one level further, into the compute layer. AWS has described a direct 1:1:1 relationship between a DynamoDB partition, its corresponding open stream shard, and the Lambda instance that processes records from that shard. Concurrency in stream processing is therefore driven by shard count, not by a setting a team tunes independently. This has a direct design consequence for partition key choice. Grouping related entities under a single partition key, such as a composite key combining a SKU and a warehouse ID, ensures every change for that logical entity lands in the same shard and gets processed in strict order by a single consumer. That avoids head-of-line blocking between unrelated entities that happen to share a table but have no business being ordered against each other. Because the mapping is 1:1 and automatic, shard count is not a tuning knob available to an engineering team; it is a direct reflection of how many partitions the table currently has.
Shard rollover: closing after up to four hours
A stream shard does not live forever. Its maximum lifespan runs up to four hours, after which the partition it serves begins publishing to a newly created shard, and the original shard is marked read-only and stamped with an EndingSequenceNumber.
A consumer walking through a shard's records will eventually hit that end marker. At that point it has to detect the EndingSequenceNumber and move to the successor shard to keep receiving records for that partition. Skipping that step throws no error and logs no warning; the consumer simply stops receiving new records, silently, which is one of the most common causes of quietly missed writes in hand-built stream consumers. Rollover is driven by time and produces exactly one successor shard, while a split (covered next) is driven by throughput and produces two child shards. Treating the two as interchangeable leads to incorrect lineage traversal and, eventually, dropped data.
A related but separate hazard sits next to rollover: disabling and re-enabling a stream on a table produces an entirely new stream ARN. Any changes that occurred during the window the stream was disabled are gone permanently, and every consumer pointed at the old ARN has to be redirected to the new one.
Shard splits: when a partition divides, its shard divides with it, and lineage becomes mandatory
When DynamoDB splits a table partition, whether because of growing data volume or rising throughput demand, the stream shard tied to that partition splits along with it. The parent shard is marked read-only, and two child shards are created, one for each of the two new partitions that resulted from the split.
That split makes lineage a hard requirement. Shards carry parent-child relationships, and a consumer has to fully process a parent shard before it starts processing either child shard. Get that order wrong, by jumping into a child shard before the parent is finished, and the consumer produces out-of-order records. That is a correctness failure, not a slow-performing edge case. The sequence number itself does not reset across the split: the SequenceNumber inside a given shard remains the correct way to order records within that shard, but establishing order across the parent and its children depends on walking the lineage graph, not on comparing sequence numbers directly.
Shard count scales with table capacity automatically. As a table's write throughput grows, the number of shards grows with it, with zero configuration required on the engineering side. That convenience hides a trap for consumers built carelessly: a consumer that enumerates the available shards once at startup and holds onto that static list will miss every new shard a later split creates. AWS's documentation on DynamoDB Streams states that the service manages shard creation and splitting transparently, placing the burden on the consumer to traverse the lineage correctly. A consumer built to handle rollover but not splits will run cleanly for months against a stable, low-growth table and then fail silently the moment that table starts scaling, which is precisely the asymmetry that makes split handling the harder and more consequential of the two lifecycle mechanics to get right.
Stream view types: record payloads and downstream effects
Enabling a stream requires choosing a view type, and that choice is made once and is immutable for the life of that stream. The view type determines what payload rides along with every stream record, and changing it later means disabling and re-enabling the stream entirely, which produces a new ARN and loses whatever was written during the gap, as described above for rollover-adjacent configuration changes.
Four view types are available, and each trades payload size against downstream capability. KEYS_ONLY sends only the partition key and sort key of the item that changed, keeping the payload smallest and suiting consumers that just need to know an item changed, such as a cache invalidator that looks up current state itself. NEW_IMAGE carries the full item as it looks after the change, and it is the most common choice for downstream synchronization where the consumer only ever wants the latest state. OLD_IMAGE carries the full item as it looked before the change, which matters for consumers building audit trails that need to know what was overwritten. NEW_AND_OLD_IMAGES carries both, producing the largest payload of the four, and it enables diff-based synchronization, conflict detection, and complete audit logging in a single record. A REMOVE event under this view type arrives with a null NewImage, so any consumer relying on it has to null-check before touching that field.
DynamoDB's schemaless design means there is no schema enforced across items in the same table, which creates a subtler problem underneath all four view types. When an application renames a field, the stream does not record that as a rename. It records the old field being removed and a new field being added, as two separate observations. A downstream relational target receiving that stream will accumulate both columns over time, with the old one going null for every record written after the rename. It is what happens when a flexible document model feeds a relational destination that expects stable columns, and it means any team syncing DynamoDB into a warehouse needs real schema-drift detection built into the consumer, not an assumption that the shape of incoming records will stay fixed.
The 24-hour retention cliff and the fan-out ceiling as hard operational constraints
Stream records live for exactly 24 hours and no longer. There is no extension available, no replay mechanism beyond that window, and no way to recover a record once it expires. A consumer that is down, throttled, or simply falling behind for longer than 24 hours has lost those writes permanently, with no recovery path inside the native stream.
The second hard constraint is fan-out. A table supports at most two simultaneous stream consumers, for example two Lambda event source mappings attached to the same stream ARN. Teams that need more than two, say an analytics pipeline, a replication target, an audit system, and a cache invalidator all running off the same change feed, cannot get there by adding more consumers directly to the stream. They need an intermediate layer, such as Kinesis Data Streams or an SNS/SQS fan-out pattern, standing between the table and that larger set of downstream consumers.
Neither constraint is a performance trade-off. The stream runs asynchronously and has no performance impact on the table itself, and stream records are encrypted at rest. Read throughput available from the stream defaults to up to twice the table's provisioned write capacity, which is generous, and Lambda-based consumers incur no additional read cost for consuming it (non-Lambda consumers pay a small per-request fee once they exceed a modest monthly free tier). Cost and throughput are not where the real limits sit. Retention and consumer count are the two boundaries that have to shape a pipeline's design from the first architecture decision, not patched in later once a team discovers them in production.
The stream's exactly-once guarantee stops at the record level and does not extend to consumer delivery
The stream itself guarantees that each record appears exactly once, in the correct order, within its shard. That guarantee stops at the boundary of the stream. What a consumer actually receives operates under at-least-once delivery: a Lambda function that crashes partway through processing a batch will have that entire batch retried, and a retry can hand the same records to the handler more than once. The gap between those two guarantees is exactly where records get duplicated in real production systems, and it is the single most consequential misunderstanding about DynamoDB Streams reliability.
An engineer who reads "exactly-once" in the documentation and stops there will build a handler that assumes it never sees the same record twice. That assumption breaks the first time a batch retry happens, producing duplicate writes in whatever downstream system is on the receiving end. The standard mitigation is to treat each record's SequenceNumber as an idempotency key: store it somewhere durable, such as a DynamoDB table with a TTL set longer than the maximum retry window, and check it before processing a record so a repeat delivery gets recognized and skipped. The stream-level exactly-once guarantee is real and documented. The at-least-once behavior a consumer actually experiences is a consequence of Lambda's own retry semantics, not a flaw in the stream.
The operational signal that ties all of this together is IteratorAge. When that metric climbs into minutes, it means the consumer is falling behind the stream's pace. The 24-hour retention wall sits behind every stream, so a rising IteratorAge is the leading indicator of impending record loss, not merely a latency annoyance to note and move past.
Routing DynamoDB changes through Kinesis Data Streams
Native DynamoDB Streams has two constraints that no amount of careful consumer engineering can work around: the 24-hour retention window and the two-consumer fan-out ceiling. Kinesis Data Streams for DynamoDB, the integration in which DynamoDB writes change records directly into a Kinesis stream, exists specifically to address both. Retention extends out to as long as one year instead of 24 hours, and Kinesis supports substantially more simultaneous consumers per shard through enhanced fan-out, well beyond the two that native DynamoDB Streams allows.
That extended capability comes with trade-offs that have to be weighed. Kinesis carries its own provisioning and cost considerations, where native DynamoDB Streams, by contrast, incurs no additional read cost for Lambda-based consumption. The decision between the two comes down to a fairly clean rule. Teams with more than two downstream consumers, teams that need a recovery window longer than 24 hours, or teams with compliance requirements demanding extended replay capability should route through Kinesis Data Streams. Teams that need strict per-item ordering, can tolerate building idempotency handling themselves, and want the simpler operational model should stay on native DynamoDB Streams.
For teams replicating DynamoDB changes into a data warehouse such as Snowflake, BigQuery, Redshift, Databricks, or ClickHouse, the operational burden compounds quickly: shard lineage has to be tracked correctly, idempotency has to be enforced at the consumer, schema drift from DynamoDB's schemaless model has to be detected before it corrupts a relational target, fan-out has to be managed if more than two consumers are in play, and the 24-hour retention wall has to be respected at all times regardless of which stream interface sits underneath. Building and operating a consumer that handles all five of those concerns correctly, on either stream interface, is substantial infrastructure work in its own right. That is the argument for treating those concerns as first-class infrastructure handled by a managed replication service, rather than as application code a team maintains indefinitely alongside its actual product.


