skip to content

How does the Kinesis Client Library distribute a stream's shards across worker instances and remember progress, and what operational problems come from that design?

level: seniorimportance: nice to knowfreq 36%

answer

  1. coordination lives in DynamoDB, not Kinesis
  2. one lease item per shard
  3. heartbeat counter, conditional steal
  4. checkpoint means at-least-once
  5. app name is the progress identity

basics

~20 s

The KCL keeps a DynamoDB lease table, one item per shard, holding the owning worker and the last checkpointed sequence number. Workers heartbeat to hold leases and steal stale ones, which makes processing at-least-once and puts a second billable table in the critical path.

solid answer

~50 s

The KCL does not use any server-side consumer-group feature — Kinesis has none. It builds coordination itself in a DynamoDB table named after your application name, with one item per shard holding `leaseKey`, `leaseOwner`, `leaseCounter` and the last checkpoint sequence number. Each worker periodically increments the counters of leases it holds; a worker that sees a counter go unchanged past the failover interval concludes the owner is dead and takes the lease over, resuming from the stored checkpoint. Workers also steal leases from each other to balance shard counts. Three consequences matter in production. Delivery is **at-least-once** — records processed after the last checkpoint are replayed on failover, so handlers must be idempotent. The lease table is a real dependency: throttle it and leases expire, workers thrash, and the stream lags. And the application name **is** the identity of the progress: reuse it in two deployments and they fight over one set of leases; change it and the new application starts from `TRIM_HORIZON` or `LATEST` with no history.

code

bash · 13 lines
bash
# The KCL's coordination state is an ordinary DynamoDB table, named after
# the application name. Inspect it to see who owns what and where they are.
aws dynamodb scan \
  --table-name clickstream-processor \
  --projection-expression "leaseKey, leaseOwner, leaseCounter, checkpoint"

# Lease churn usually shows up as throttling on this table.
aws cloudwatch get-metric-statistics \
  --namespace AWS/DynamoDB \
  --metric-name ThrottledRequests \
  --dimensions Name=TableName,Value=clickstream-processor \
  --start-time 2025-01-01T00:00:00Z --end-time 2025-01-01T01:00:00Z \
  --period 300 --statistics Sum

go deeper

for a junior

Know that Kinesis itself does not remember where a consumer got to — the client library does, by storing checkpoints in a DynamoDB table named after your application.

for a middle

Explain the lease item's fields and the heartbeat-and-steal mechanism, and connect checkpoint frequency to the duplicate window a failover produces.

for a senior

Diagnose lease churn: recognise DynamoDB throttling on the lease table as a cause of consumer lag, and know why child shards idle until their parents reach SHARD_END after a reshard.

for a principal

Own the invariants around this state: application name as production identity with a rename policy, lease-table capacity and cost as part of the consumer's budget, and idempotency as a design requirement rather than a code review note.

## Kinesis has no consumer groups, so the client builds one On the service side, Kinesis Data Streams offers shards, iterators, and `GetRecords`/`SubscribeToShard`. There is no server-side notion of a group of workers, no assignment protocol, and no stored progress. Everything about "which of my six machines reads which shard, and where did we get to" is a **client-side** construct, implemented by the Kinesis Client Library (KCL) — and, in the same shape, by the Lambda and Flink integrations. ## The lease table On first start, the KCL creates a DynamoDB table named after the application name you configured. It holds one item per shard: - `leaseKey` — the shard id, and the table's partition key - `leaseOwner` — the worker identifier currently processing it - `leaseCounter` — a monotonically incremented heartbeat - `checkpoint` — the sequence number of the last record the owner declared done (or the sentinels `TRIM_HORIZON`, `LATEST`, `SHARD_END`) - parent shard ids, so resharding can be handled correctly Holding a lease means periodically incrementing `leaseCounter` with a conditional write. A worker takes a lease by conditionally writing its own id when the counter it last observed is unchanged after the failover interval — the classic conditional-write lease. Because the write is conditional on the counter, two workers cannot both win. Balancing rides on the same mechanism: a worker holding fewer than its fair share of leases steals one from the worker holding the most, again by conditional write. Start a second instance of a six-shard application and it converges to three leases each without any coordinator being involved. ## Checkpointing and delivery semantics Your record processor calls `checkpoint()` when it has durably handled records up to a point; the KCL writes that sequence number into the lease item. On failover, the new owner reads the checkpoint and resumes there. That design is squarely **at-least-once**. Records processed but not yet checkpointed when a worker dies are delivered again to whoever picks up the lease. Checkpoint after every record and you minimise duplicates but write to DynamoDB constantly; checkpoint every N records or every few seconds and you trade a wider replay window for far fewer writes. Either way the processing must be idempotent — deduplicate on a business key, or make the write itself idempotent (conditional put, upsert by id). ## Resharding A `SplitShard` or `UpdateShardCount` closes parent shards and opens children. The KCL discovers this, writes lease items for the children recording their parents, and enforces a rule: a child shard is not started until every parent has been read to `SHARD_END` and checkpointed. Since a given partition key lived in the parent before the split and in exactly one child after it, that ordering rule is what preserves per-key ordering across a reshard. It also explains a common observation — after resharding, some new shards sit idle for a while; they are waiting on their parents, not broken. ## The operational problems **The lease table is a dependency you did not ask for.** It is billed, it can be throttled, and it can be deleted by an over-eager cleanup script. Under-provisioned capacity means lease renewals fail, workers believe each other dead, leases churn between them, and throughput collapses even though the stream itself is healthy. Symptoms: rising `GetRecords.IteratorAgeMilliseconds`, DynamoDB throttled-request metrics on that table, and log lines about lost or stolen leases. Provision it for the write rate your checkpoint interval implies, or use on-demand billing for it. **The application name is the identity of progress.** Two deployments (staging and production, blue and green) sharing an application name share one lease table and will each think the other's workers are theirs — shards get processed by the wrong environment. Conversely, renaming the application on a deploy creates a fresh table and the new one starts at whatever initial position you configured. That is either a full reprocess from `TRIM_HORIZON` or a silent gap from `LATEST`, and both have caused real incidents. **Parallelism is capped by shard count.** One lease per shard means at most one worker actively processes a shard. Run twelve instances against a six-shard stream and six sit idle. Scaling the consumer means scaling the stream, which is a slower and more consequential operation than adding pods. **Version differences.** KCL 2.x retrieves via enhanced fan-out by default and can be configured for polling; KCL 1.x polls `GetRecords`. If you migrate and see a new per-consumer-shard-hour line on the bill, that default is why. ## What to monitor Per-shard `GetRecords.IteratorAgeMilliseconds` (the primary lag signal), DynamoDB throttling and consumed capacity on the lease table, and the KCL's own lease-take and lease-loss logs. A healthy application shows stable lease ownership; leases moving constantly between workers is the fingerprint of a lease-table or timeout misconfiguration, not of a slow stream.

  • How often should a record processor checkpoint, and what does the choice trade off?
    Checkpointing after every record minimises replay on failover but writes to DynamoDB at the record rate, which is expensive and easy to throttle. Checkpointing every few seconds or every N records collapses those writes at the cost of a wider duplicate window. Pick the interval from how much reprocessing your idempotent handler can absorb cheaply, then provision the lease table for that write rate.
  • A new deployment of the consumer suddenly reprocesses days of data. What is the most likely cause?
    The application name changed, so the KCL created a fresh lease table with no checkpoints and started from the configured initial position — `TRIM_HORIZON` replays everything still retained. The mirror-image failure is starting fresh at `LATEST`, which silently skips everything written before the deploy. Treat the application name as production state and never let it drift with a rename or an environment-variable typo.
  • You run twelve consumer instances against a six-shard stream and throughput does not improve. Why?
    Because the KCL holds one lease per shard and only the lease owner processes it, so six instances own a shard each and the other six idle. Consumer parallelism is bounded by shard count. To go faster you either add shards, or move slow work off the record-processing thread into a bounded worker pool so each lease holder drains its shard faster.

saying these in an interview costs you the question

  • Assumes Kinesis stores consumer progress server-side
  • Says the KCL gives exactly-once delivery
  • Ignores the lease table's capacity and billing
  • Thinks more consumer instances than shards adds throughput
  • Treats the application name as a cosmetic label

context