Serverless and Function-as-a-Service Architecture Questions
Building applications on managed, event-triggered compute: functions-as-a-service (AWS Lambda, Azure Functions, Google Cloud Functions, Cloudflare Workers) and serverless containers. Covers the invocation lifecycle and cold starts (init vs handler work, provisioned concurrency, packaging, layers and container images), statelessness and externalizing state, execution limits (timeouts, memory, payload and /tmp size), concurrency and scaling behavior (account limits, burst scaling, protecting downstream databases with connection proxies and throttling), event sources and trigger semantics (at-least-once delivery, retries, idempotent handlers, dead-letter handling), composing functions with managed services and workflow orchestrators (Step Functions and equivalents), and serverless-specific observability, security (per-function IAM, secrets) and pay-per-invocation cost modeling. Includes when serverless fits versus containers or VMs, and vendor lock-in trade-offs. General compute selection, generic messaging patterns, and general idempotency theory are covered by their own topics.
Design a serverless, event-driven pipeline that ingests telemetry at 100k events/second, does lightweight enrichment and aggregation, and writes results to analytic storage. Cover event ingestion and buffering, processing concurrency, idempotency, error handling, storage choice, observability, and the cost profile of your design.
Sample Answer
Direct answer
At a sustained 100,000 events per second I would build this as Kinesis Data Streams (provisioned shards) feeding a Lambda consumer that enriches and pre-aggregates in batches, with a per-shard checkpoint committed atomically alongside the aggregates, plus a Firehose copy of the raw stream into S3 for replay. The two decisions that matter most are both about batching: producers pack many events into each stream record, and Lambda is invoked once per batch of about 1,600 events rather than once per event. Invoking once per event would cost about $52,560 a month in request fees alone. The batched design runs at about $10,900 a month, and the raw archive (not the compute) is the biggest line item.
Quick glossary for the rest of the answer:
- Serverless / FaaS (functions as a service): you upload a function, the cloud runs as many copies as traffic needs and bills per invocation and per millisecond. AWS Lambda is the example here.
- Kinesis Data Streams (KDS): AWS's managed, ordered, replayable log (similar in spirit to Kafka, Apache Kafka: the most common self-hosted alternative doing the same ordered-log job). It is split into shards; each shard accepts up to 1 MB/s or 1,000 records/s of writes.
- Event source mapping (ESM): the Lambda-managed poller that reads a shard and invokes your function with a batch of records.
- At-least-once delivery: every record is delivered, but some may be delivered twice (after a retry). The pipeline must tolerate duplicates.
- Idempotent: processing the same input twice leaves the same result as processing it once.
- Amazon Data Firehose: a managed service that buffers a stream and writes large files to S3. S3 is object storage; Parquet is a columnar file format; Athena is a serverless SQL engine that queries files in S3.
- DynamoDB: AWS's managed key-value database.
Requirements I am designing to
| Assumption | Value | Why it matters |
|---|---|---|
| Throughput | 100,000 events/s sustained | Drives shard count and cost |
| Event size | about 1 KB JSON | 100 MB/s, 262,800 GB per 730-hour month |
| Freshness | aggregates queryable within about 3-5 minutes | Allows 2 s batching windows and a short grace period |
| Aggregation | per-minute counts and sums by (device group, metric) | Keeps aggregate state small |
| Correctness | aggregates must not double count | Delivery is at-least-once, so this needs design |
Architecture
flowchart LR
P[Devices and gateways] -->|packed 24 KB records| K[Kinesis stream 125 shards]
K --> L[Lambda consumer: enrich and pre-aggregate]
L -->|one transaction per batch| D[(DynamoDB: per-shard minute aggregates and checkpoint)]
L -.->|poison records| Q[SQS on-failure queue]
D --> R[Rollup job every minute]
R --> S3A[(S3 aggregates as Parquet, queried by Athena)]
K --> F[Firehose]
F --> S3R[(S3 raw archive for replay)]
1. Ingestion and buffering
- Shard count. 100 MB/s divided by 1 MB/s per shard is 100 shards by bytes. With 24 events packed per record, the record rate is 100,000 / 24 = 4,167 records/s, far below the 100,000 records/s limit of 100 shards. Bytes bind, so I provision 125 shards (25% headroom for skew and bursts).
- Why pack events. Unpacked, 100,000 records/s would need 100 shards just for the record-count limit, and every downstream fee charged per record (Kinesis PUT payload units, which bill each record in 25 KB chunks; Firehose's 5 KB billing increment) would be charged 100,000 times a second. Packing is done by the gateway or by the Kinesis Producer Library's (KPL, an AWS client library that batches and compresses records before sending them) aggregation feature.
- Partition key. Device ID, hashed, so load spreads evenly across shards. A partition key concentrated on a few tenants creates a hot shard (one shard at its 1 MB/s ceiling while others idle), which shows up as write throttling.
- Provisioned vs on-demand. On-demand mode (no shard planning, billed per GB) costs about $42,048/month at this volume versus about $1,522 provisioned (arithmetic below). For flat, known 100k/s load, provisioned wins by more than 20x. I would flip to on-demand only for a new stream whose traffic I cannot yet predict.
- Buffering. The stream itself is the buffer: default retention is 24 hours, so a consumer outage of a few hours loses nothing; the backlog is replayed when the consumer recovers.
2. Processing concurrency
- The Kinesis ESM runs one concurrent batch per shard by default.
ParallelizationFactor(1 to 10) raises that, but it breaks per-shard ordering (ordering is then only per partition key), and my deduplication below relies on per-shard order. So I keep the factor at 1 and scale by adding shards instead. - Maximum concurrency is therefore 125 execution environments, well inside Lambda's default regional quota of 1,000 concurrent executions. I set reserved concurrency (a per-function concurrency allotment that is both a guaranteed floor and a hard ceiling) of about 150 on this function so it can never starve other functions, and so nothing else can starve it.
- Batching settings:
BatchSize10,000 (the maximum),MaximumBatchingWindowInSeconds2. Each shard receives 100,000 / 125 = 800 events/s, so each 2 s batch carries about 1,600 events (about 67 packed records, about 1.6 MB, safely under the 6 MB invocation payload limit even after base64 encoding inflates it by a third). - That gives 125 / 2 = 62.5 invocations/s. If a batch takes about 120 ms, average concurrency is 62.5 x 0.12 = 7.5 environments. Cold starts (the extra latency when a new execution environment boots) are irrelevant here: environments stay warm under constant load, and a few hundred ms of extra delay on a 2 s batch changes nothing.
- Enrichment (joining each event to device metadata) uses a reference table loaded into memory at initialization and refreshed every few minutes, never a per-event database call. 100,000 lookups/s against a database would itself be a larger system than the pipeline.
3. Idempotency: exactly-once aggregates on at-least-once delivery
Kinesis plus Lambda is at-least-once. The ESM retries a whole batch if the function errors or times out, so if the function wrote aggregates and then timed out before returning, the retry would add the same events again. For aggregates (counters), duplicates silently inflate numbers.
The fix: store a per-shard checkpoint (the highest sequence number already counted) in the same atomic write as the aggregate increments. Each Kinesis record has a sequence number that increases within a shard. On every invocation the function:
- Reads the checkpoint for its shard (cached in memory, re-read if the conditional write in step 3, a database write that only takes effect if the value has not changed since it was last read, fails).
- Drops records whose sequence number is at or below the checkpoint (already counted).
- Commits, in one DynamoDB
TransactWriteItemscall (an all-or-nothing multi-item write): increment the (shard, minute) aggregate items, and set the checkpoint to the batch's last sequence number, with a condition that the checkpoint still equals the value just read.
Either both the counts and the checkpoint land or neither does, so a retry of any subset of the batch finds its records already covered and adds nothing.
A small simulation proves it. It models 4 shards of 5,000 records, with a 10% chance each attempt crashes before writing and a 10% chance it crashes after writing, and halves the batch after each failure (a simplified model of BisectBatchOnFunctionError):
import random
from collections import Counter
def make_stream(num_shards=4, per_shard=5000, seed=7):
rng = random.Random(seed)
shards = {}
for s in range(num_shards):
recs = []
for seq in range(1, per_shard + 1):
minute = (seq - 1) * 60 // per_shard # events spread over 60 minutes
group = rng.choice(["sensor-a", "sensor-b", "sensor-c"])
recs.append((seq, minute, group))
shards[s] = recs
return shards
class Store:
"""Stands in for the database: aggregate counters plus a per-shard checkpoint.
commit() is atomic, like one TransactWriteItems call."""
def __init__(self):
self.agg = Counter()
self.hwm = {}
def commit(self, shard, deltas, new_hwm):
self.agg.update(deltas)
if new_hwm is not None:
self.hwm[shard] = new_hwm
def handler(store, shard, batch, use_checkpoint):
last_done = store.hwm.get(shard, 0) if use_checkpoint else 0
fresh = [r for r in batch if r[0] > last_done]
deltas = Counter((shard, minute, group) for _, minute, group in fresh)
new_hwm = max(last_done, batch[-1][0]) if use_checkpoint else None # never move backwards
store.commit(shard, deltas, new_hwm)
def run(use_checkpoint, batch_size=500, p_fail_before=0.10, p_fail_after=0.10, seed=42):
shards = make_stream()
rng = random.Random(seed)
store = Store()
for shard, recs in shards.items():
pos, size = 0, batch_size
while pos < len(recs):
batch = recs[pos:pos + size]
roll = rng.random()
if roll < p_fail_before: # crash before the write: nothing committed
size = max(1, len(batch) // 2) # simplified bisect-on-error
continue
handler(store, shard, batch, use_checkpoint)
if roll < p_fail_before + p_fail_after: # write committed, then timeout: batch retried
size = max(1, len(batch) // 2)
continue
pos += len(batch)
size = batch_size
truth = Counter((s, m, g) for s, recs in shards.items() for _, m, g in recs)
return truth, store.agg
for mode in (False, True):
truth, got = run(use_checkpoint=mode)
label = "with per-shard checkpoint" if mode else "naive (no checkpoint) "
print(f"{label}: true events={sum(truth.values())}, counted={sum(got.values())}, "
f"every (shard, minute, group) cell exact={truth == got}")
Output:
naive (no checkpoint) : true events=20000, counted=22250, every (shard, minute, group) cell exact=False
with per-shard checkpoint: true events=20000, counted=20000, every (shard, minute, group) cell exact=True
The naive version over-counts by 11.25%. The max(...) on the checkpoint line is load-bearing: my first version set the checkpoint to the retried half-batch's last sequence number, which moved it backwards after a bisect and still double counted (21,125). A checkpoint must only ever move forward.
Two things this does not cover, stated so nobody assumes otherwise: duplicates created by the producer (a device retrying a PUT gets a new sequence number) need an event ID and a dedupe step, and the raw archive in S3 is at-least-once, so queries over raw events should deduplicate by event ID.
4. Error handling
| Failure | Handling |
|---|---|
| Transient (throttle, timeout, database conflict) | Function raises; ESM retries. MaximumRetryAttempts 5, MaximumRecordAgeInSeconds 3,600 so a stuck shard cannot block for the whole 24 h retention |
| One malformed record | Validate and skip it inside the handler, writing it to an error log with its sequence number. It must not fail the batch |
| A record that crashes the code | BisectBatchOnFunctionError splits the batch to isolate it; ReportBatchItemFailures (partial batch response) lets the function say "everything before sequence X succeeded" so good records are not reprocessed |
| Retries exhausted | ESM on-failure destination: an SQS (Amazon Simple Queue Service) queue that receives the batch's shard and sequence range (metadata, not the records), from which an operator replays the records out of the stream |
| Late events | The minute rollup runs at minute + 3 and again at minute + 10, overwriting that minute's output, so late arrivals within 10 minutes are counted |
The DLQ (dead-letter queue) idea here is the on-failure destination: a place failed work goes so it is not lost and does not block the shard.
5. Storage choice
- Aggregates: DynamoDB holds the live per-shard, per-minute partial counts (small, write-heavy, needs atomic conditional writes). A rollup (a job that combines many small partial records into one summarized one) Lambda run every minute by EventBridge Scheduler (AWS's managed cron) sums the 125 shard items for a closed minute and writes one Parquet object per minute to S3 with a deterministic key such as
aggregates/dt=2026-09-27/hh=14/mm=05.parquet. Rewriting the same key is naturally idempotent. Analysts query with Athena. A small daily compaction merges per-minute files into hourly ones to avoid a small-file problem (many tiny files make each file's fixed per-file read overhead dominate, slowing every downstream query). - Raw events: Firehose reads directly from the stream and writes compressed files to S3, so enrichment or aggregation logic can be replayed after a bug fix.
- What would change it: if dashboards need sub-second freshness or high-concurrency interactive queries, I would write aggregates to a real-time OLAP (online analytical processing) database such as ClickHouse or Apache Druid instead of S3 plus Athena.
6. Observability
- Consumer lag: Lambda's
IteratorAgemetric (how old the last record in each batch was when the batch was sent to the function). This is the single most important alarm: rising iterator age means processing is falling behind. Alarm if it exceeds 60 s for 5 minutes. - Stream health: Kinesis
WriteProvisionedThroughputExceeded(a hot or under-provisioned shard) andGetRecords.IteratorAgeMilliseconds. - Function health:
Errors,Throttles,Duration(p99, meaning the 99th percentile),ConcurrentExecutions. - Business correctness: a reconciliation job (one that independently checks two derived numbers against each other and flags any mismatch) comparing hourly raw-archive event counts with aggregated counts, alarming on a difference above 0.1%. This catches silent double counting or silent drops that no infrastructure metric will show.
- Structured logs with shard ID and sequence range per batch, so any aggregate can be traced to the records behind it.
7. Cost profile
List prices for us-east-1 as published on the AWS pricing pages; the per-event CPU cost is an assumption to be measured in a load test. DynamoDB is billed in write request units (WRU): under on-demand pricing a normal write costs 1 WRU per KB, and (as the pitfall below explains) a transactional write costs double that per KB:
import math
# us-east-1 list prices (verify on the pricing pages before relying on them)
LAMBDA_GBS = 0.0000166667 # $ per GB-second, x86, first tier
LAMBDA_REQ = 0.20 / 1e6 # $ per request
KDS_SHARD_HR = 0.015 # provisioned shard-hour
KDS_PUT_UNIT = 0.014 / 1e6 # per 25 KB PUT payload unit
KDS_OD_IN, KDS_OD_OUT = 0.08, 0.04 # on-demand $/GB ingested, retrieved
FH_GB = 0.029 # Firehose ingestion $/GB, billed in 5 KB increments
DDB_WRU = 0.625 / 1e6 # on-demand write request unit
SEC_MONTH = 730 * 3600 # AWS bills a month as 730 hours
events_s, event_kb = 100_000, 1.0
events_per_record = 24 # producers pack ~24 KB records
shards = 125 # 100 MB/s needs 100 shards at 1 MB/s each, plus 25% headroom
records_s = events_s / events_per_record
mb_s = events_s * event_kb / 1000
gb_month = mb_s / 1000 * SEC_MONTH
kds_prov = shards * 730 * KDS_SHARD_HR + records_s * SEC_MONTH * KDS_PUT_UNIT
kds_od = gb_month * (KDS_OD_IN + 2 * KDS_OD_OUT) # two consumers: Lambda and Firehose
window_s, mem_gb = 2, 1.0
inv_s = shards / window_s
events_per_inv = events_s / inv_s
duration_s = 0.040 + events_per_inv * 0.00005 # ASSUMED: 40 ms fixed + 0.05 ms per event
lam = inv_s * SEC_MONTH * (LAMBDA_REQ + duration_s * mem_gb * LAMBDA_GBS)
per_event_invoke_requests = events_s * SEC_MONTH * LAMBDA_REQ
wru_per_inv = 2 * 1 + 2 * 4 * 1.05 # txn: 1 KB checkpoint + ~1.05 x 4 KB aggregate item, 2 WRU per KB
ddb = inv_s * SEC_MONTH * wru_per_inv * DDB_WRU
fh_billed_kb = math.ceil(events_per_record * event_kb / 5) * 5
fh = records_s * fh_billed_kb / 1e6 * SEC_MONTH * FH_GB
print(f"records/s={records_s:.0f} data={mb_s:.0f} MB/s GB/month={gb_month:,.0f}")
print(f"Kinesis provisioned={kds_prov:,.0f} on-demand={kds_od:,.0f}")
print(f"Lambda: {inv_s:.1f} inv/s, {events_per_inv:.0f} events/inv, {duration_s*1000:.0f} ms, "
f"avg concurrency={inv_s*duration_s:.1f}, cost={lam:,.0f}")
print(f"Lambda request fee alone if invoked once per event={per_event_invoke_requests:,.0f}")
print(f"DynamoDB aggregates={ddb:,.0f}")
print(f"Firehose raw archive (billed {fh_billed_kb} KB/record)={fh:,.0f}")
print(f"TOTAL (provisioned Kinesis)={kds_prov + lam + ddb + fh:,.0f} $/month")
Output (dollars per month):
records/s=4167 data=100 MB/s GB/month=262,800
Kinesis provisioned=1,522 on-demand=42,048
Lambda: 62.5 inv/s, 1600 events/inv, 120 ms, avg concurrency=7.5, cost=361
Lambda request fee alone if invoked once per event=52,560
DynamoDB aggregates=1,068
Firehose raw archive (billed 25 KB/record)=7,939
TOTAL (provisioned Kinesis)=10,890 $/month
The ~4 KB aggregate-item size is an assumption for this example: it is enough room for counts and sums across roughly a few dozen (device group, metric) key combinations packed into one item; size it against your own key cardinality. The ~1.05 multiplier on that item accounts for the occasional batch that straddles a minute boundary and touches two minute items. S3 storage and Athena query charges are excluded because they depend on retention and query volume.
What the numbers say:
- Compute is the cheap part ($361). Lambda cost is dominated by how you batch, not by how fast the code is.
- The raw archive is 73% of the bill. Its value (replay, ad hoc analysis) should be confirmed with the data consumers. If a 7-day replay window is enough, extending the stream's retention instead of archiving everything could remove most of it.
- Firehose bills each record rounded up to 5 KB, so 1 KB unpacked records would be billed at 5x their size. Packing records to just under a 5 KB multiple matters.
Trade-offs and pitfalls
- Is serverless right at this volume? The load is flat and high, which is where serverless is weakest on unit price. It still wins here because the managed pieces (stream, poller, retries, scaling to 125 parallel consumers, no servers to patch) cost less in engineering time than running and patching your own Kafka and stream-processing cluster, while the Lambda line itself is only a few hundred dollars. If enrichment grew heavy, say 2 ms of CPU per event instead of 0.05 ms, two things break. First, one consumer per shard can then process at most 1 / 0.002 = 500 events/s, but each shard receives 800, so the stream falls steadily behind unless you add shards: 100,000 events/s / 500 events/s per shard needs at least 200 shards by that arithmetic (up from 125), before headroom. Second, compute becomes 100,000 x 0.002 = 200 GB-seconds every second at 1 GB, which is 200 x 2,628,000 s/month x $0.0000166667/GB-s, about $8,760 a month in Lambda duration alone. At that point a long-running container consumer (for example Apache Flink, an open-source stream-processing engine that keeps one continuously running process instead of Lambda's start-stop-per-batch model, on a managed service) is the better call, because steady, CPU-heavy work is exactly where per-millisecond billing loses.
- Per-event invocation is the classic mistake: 262.8 billion requests a month.
- ParallelizationFactor as the first scaling lever silently breaks the per-shard ordering the checkpoint relies on. Scale with shards.
- Aggregating in the function's memory across invocations is unsafe: environments are recycled at any time and two environments can serve the same shard over time. State lives in the database or in Lambda's managed tumbling-window state (a built-in feature that carries a small amount of state between consecutive batches from the same shard, for windowed aggregation, which AWS documents as at-least-once and capped at 1 MB per shard, so it does not remove the double-count problem).
- Transactions are not free: each item written in a DynamoDB transaction costs 2 write request units (WRU) per KB instead of 1. That is why there is one transaction per batch, not per event.
Explain stateless versus stateful function design in a serverless architecture. When should a function be fully stateless, and when is holding state (session context, aggregated results, a warm cache) actually justified? Where does that state live once it's outside the function, and what latency or consistency trade-offs come with each option?
Sample Answer
Direct answer
A stateless function keeps nothing it depends on between invocations: everything it needs arrives in the event or is read from an external store, and everything worth keeping is written back out. That is the default in serverless, because the platform runs many copies in parallel, routes each request to whichever copy is free, and destroys copies whenever it likes. Holding state inside the function is justified only as a disposable optimization: a warm cache of read-mostly data or a reused connection, where losing it costs a slower request but never a wrong answer. Anything that must be correct or durable (sessions, running totals, workflow progress) lives outside the function, and you pick the store by the latency and consistency you need.
Why "in-memory state" breaks in FaaS
Three properties of the platform, each enough on its own:
- Many copies. At 50 concurrent requests there are about 50 separate execution environments (sandboxes), each with its own memory. A counter in one is invisible to the others.
- No affinity. A user's second request can land on a different environment than their first, so an in-memory session is simply missing.
- No lifetime guarantee. Environments are frozen when idle, recycled every few hours even under load, and reset after a crash. Memory contents vanish without warning.
The simulation below routes 1,000 click events to environments the way a platform might (any environment, new ones added under load), and has the platform recycle one environment part-way through. Each environment keeps a local counter; the "external" counter stands in for an atomic counter in a database (one whose increments cannot partially overlap: two simultaneous +1s always land as +2, never silently lose one).
import random
random.seed(7)
EVENTS = 1000
envs = {} # env_id -> in-memory click count (what a "stateful" function would keep)
external_total = 0 # stands in for an atomic counter in an external store
next_env = 0
def pick_env():
global next_env
# the platform routes to any idle environment, and adds new ones under load
if not envs or random.random() < 0.01:
envs[next_env] = 0
next_env += 1
return random.choice(list(envs))
for i in range(EVENTS):
e = pick_env()
envs[e] += 1 # in-memory aggregation
external_total += 1 # externalized aggregation
if i == 600: # the platform recycles one environment mid-stream
victim = min(envs)
lost = envs.pop(victim)
print("environments created:", next_env)
print("any single environment's count:", envs[min(envs)])
print("sum of surviving in-memory counts:", sum(envs.values()))
print("counts lost when env", victim, "was recycled:", lost)
print("external counter:", external_total)
Output:
environments created: 15
any single environment's count: 163
sum of surviving in-memory counts: 792
counts lost when env 0 was recycled: 208
external counter: 1000
No single environment knows the true total (163 is one environment's view), the in-memory sum is short by the 208 clicks held in the recycled environment, and only the externalized counter is right.
When holding state is justified
| State kept in the function | Justified? | Condition |
|---|---|---|
| Reused database or HTTP client connection | Yes | Reconnect if it has gone stale after a freeze. |
| Cache of read-mostly data (config, feature flags, meaning switches that turn a feature on or off without a deploy, reference tables, a loaded model) | Yes | Has a TTL (time-to-live, an expiry) or version check; a miss just means one slower request. |
| Per-user session context | No | Another environment will serve the next request. Put it in an external store keyed by session ID. |
| Aggregated results (counters, running sums, batches being built) | No | Copies split the data and recycling loses it. Aggregate in the store, or stream events to something built to aggregate. |
| Progress of a multi-step job | No | A crash or timeout loses it. Use a workflow orchestrator (for example AWS Step Functions, a managed state machine that records each step's result) or a job table. |
The test: if this state disappeared right now, would the next request be slower, or wrong? Slower is acceptable for an in-function cache. Wrong means it belongs outside.
Where externalized state lives, and what each option costs you
| Store | Typical use | Latency (order of magnitude) | Consistency to know about |
|---|---|---|---|
| Key-value database (Amazon DynamoDB, Azure Cosmos DB, Firestore) | Sessions, idempotency records (a stored marker of which request IDs have already been processed, so a retry can be detected and skipped), counters via atomic update | Single-digit milliseconds within the Region | DynamoDB reads are eventually consistent by default (a read right after a write can be stale); request a strongly consistent read when you must see your own write. Atomic ADD/conditional updates let many copies update one counter without losing increments. |
| In-memory cache (Redis or Memcached, as a managed service such as Amazon ElastiCache) | Hot session data, rate-limit counters, shared cache | Sub-millisecond to low milliseconds, network hop included | Replication to replicas is asynchronous, so a failover can lose recent writes. Treat it as a cache unless you've configured and tested durability. On AWS it runs inside a VPC (Virtual Private Cloud, a private network), so the function must join that network. |
| Object storage (Amazon S3) | Large blobs, intermediate files, batch outputs | Tens of milliseconds per request | S3 gives strong read-after-write consistency for objects, but it is not built for small, frequent updates to the same key. |
| Relational database (PostgreSQL, MySQL) | Transactional state across entities | Low milliseconds, plus connection set-up | Strong consistency and transactions; the risk is connection count, since each concurrent environment wants its own connection. Use a connection proxy (such as Amazon RDS Proxy) or cap function concurrency. |
| Workflow orchestrator (Step Functions, Azure Durable Functions) | Multi-step progress, retries, waiting for humans | Per-step overhead, well above a single function call | The orchestrator owns the state machine, so a failed step resumes rather than restarting from zero. |
Worked example: a click-counting API
Requirement: count clicks per article per minute, 2,000 clicks per second at peak, dashboard reads the counts. Assume a lightweight handler, about 15 ms per click (just an increment call): by Little's Law (the number of requests in flight equals the arrival rate times the time each spends being handled), 2,000/s x 0.015s = 30, which is where the "~30 environments in flight" below comes from.
- Wrong: each function keeps
counts[article] += 1in memory and writes it out "every minute". With ~30 environments in flight (the arithmetic above), 30 partial counts race to overwrite each other, and a recycled environment loses its last minute. - Right, simple: each click performs an atomic increment on a key such as
article#42#2026-09-27T10:15in DynamoDB. 2,000 writes per second is fine for the table as a whole, but a single viral article sends all its writes to one key; spread that hot key (a single partition key taking disproportionate traffic) across N sub-keys (#shard0to#shardN) and sum them on read. - Right, at larger scale: put clicks on a stream (a managed, ordered, replayable log of events that many independent readers can consume, such as Amazon Kinesis or Apache Kafka) instead of writing straight to the table. A stream is split into shards, in this context a partition of the stream that one reader handles at a time (a different sense of "shard" from the DynamoDB sub-keys just above, which split a hot row, not a stream); let one consumer, a process reading a shard in order, aggregate clicks per minute and write the result once. Aggregation in a stream consumer is safe because the stream, not the function's memory, is the source of truth: if a consumer crashes partway through a minute, a new one picks up the shard from the last committed position and can be replayed, meaning it re-reads events the stream already has instead of losing them.
- The justified warm state: each function keeps a 60-second in-memory cache of the article metadata it needs to validate clicks. Losing it costs one extra read.
Trade-offs and pitfalls
- Latency versus correctness. Every external read or write adds a network hop. Cache reads locally when stale data is acceptable; never cache anything you then use to make an authoritative decision (balances, inventory).
- Globals leak across requests. Module-level variables survive between invocations in the same environment, so per-request data stored there can bleed into the next user's request. Keep request data local.
/tmpis memory-adjacent state too. Files there survive between warm invocations and vanish on recycle. Same rule: cache only.- Retries make external writes happen twice. Event sources deliver at least once, so a retried invocation can double-increment. Where exactness matters, record a per-event ID with a conditional write (a database write that only succeeds if a stated condition holds, here "this event ID is not already recorded") and skip duplicates.
How do FaaS platforms scale as traffic increases rapidly? Using AWS Lambda as an example, explain how new execution environments get provisioned, what account or regional concurrency limits exist, and what reserved versus provisioned concurrency each do. What does this mean for a latency-sensitive endpoint under a sudden burst of traffic?
Sample Answer
Direct answer
A FaaS platform (functions-as-a-service: you upload a function and the provider runs it once per event) scales by creating more execution environments. An execution environment is a small isolated sandbox with your runtime and code loaded, and it handles one request at a time. So 600 requests in flight at the same moment need 600 environments. When a request arrives and no idle environment exists, Lambda builds a new one first. That is a cold start. Growth is bounded by two different limits:
- a per-function scaling rate: how fast new environments can be created
- a Regional account concurrency quota: how many can exist at once, shared by every function in the account
Reserved concurrency gives a function its own slice of that quota and caps it there. Provisioned concurrency keeps environments already initialized, so requests that land on them skip the cold start. For a latency-sensitive endpoint, the cold starts and throttles all arrive in the same first seconds of a burst. The fix is to pre-warm the endpoint with provisioned concurrency sized to the burst and make sure the quota has headroom before the burst starts.
How a new environment gets provisioned
- Routing. A request arrives. If an idle, already-initialized ("warm") environment exists, Lambda sends the request there. No setup cost.
- Init phase (the cold start). If there is no idle environment, Lambda creates one: it loads your code package, starts the language runtime and any extensions, then runs your top-level code. That is everything outside the handler: imports, SDK clients, configuration loading.
- Invoke phase. The handler runs.
- Freeze and reuse. After the response, the environment is frozen and kept for later requests. How long it is kept is unspecified. It can be reclaimed at any time, so you cannot rely on it staying warm.
Concurrency means the number of requests being processed at the same instant, and it is also the number of execution environments Lambda needs at that instant. You can compute it with Little's law: concurrency = arrival rate (RPS) x the average time each request spends in the system (duration in seconds). That gives you environments directly; there is no hard per-environment throughput ceiling. AWS's own conceptual scaling guide warns against exactly the framing that sounds intuitive here: "it's incorrect to say that each Lambda execution environment can handle only a maximum of 10 requests per second. Instead of observing the load on any individual execution environment, Lambda only considers overall concurrency and overall requests per second when calculating your quotas." A single warm environment reused back to back for a fast function easily clears 10 requests per second in practice (a 20 ms handler run continuously on one environment serves 50 requests per second, not 10).
What Lambda actually enforces alongside the concurrency quota is a separate REQUEST-RATE quota: RPS must stay at or below 10 times whichever concurrency quota is in force (10,000 RPS against the default 1,000-unit account quota; 10 times a function's reserved concurrency if it has one; 10 times its provisioned concurrency for the units running on it). This is a policy ceiling tied to the quota number, not a physical property of any one environment. For a function whose average duration is at or above 100 ms, Little's law alone already keeps the request rate inside that 10x ceiling, so it never binds. It matters specifically for the sub-100-ms case: a 20 ms function serving 30,000 RPS needs a concurrency, and so an environment count, of only 600, comfortably under a 1,000-unit account quota, but its request rate of 30,000 is three times the 10,000 RPS the default quota allows, so it throttles on the rate ceiling even though the environment count looks fine.
environments (concurrency)=RPS×durationswith a second, independent check for functions under about 100 ms average duration:
RPS≤10×concurrency quota (account default, or the function’s reserved/provisioned concurrency)The limits
| Limit | Scope | Default | What happens when you hit it |
|---|---|---|---|
| Concurrency quota | Account, per Region, shared by all functions | 1,000 (new accounts start lower); can be raised to tens of thousands through Service Quotas (the AWS console and API where you request a higher account limit) | Requests beyond it get a throttling error (HTTP 429) |
| Concurrency scaling rate | Per function, per Region | 1,000 new environments per 10 seconds, refilled continuously (picture a token bucket that starts full at 1,000: creating an environment spends one token, and tokens refill steadily, but the bucket cannot hold more than 1,000 at once); unused capacity does not accumulate beyond that 1,000-token cap | Requests arriving faster than environments can be created are throttled (429) |
| Unreserved floor | Account | 100 units always stay unreserved | You cannot reserve or provision more than (quota minus 100) |
The effect of a throttle depends on who called the function:
- Synchronous callers (API Gateway, a direct invoke) get the error immediately. Nothing retries for you.
- Asynchronous sources (S3 notifications, SNS, Amazon Simple Notification Service, a managed pub/sub messaging service) go through Lambda's internal queue, which retries throttled events for up to 6 hours by default.
- Queue pollers (Amazon SQS, Simple Queue Service) leave the messages in the queue to be retried.
Reserved versus provisioned concurrency
| Reserved concurrency | Provisioned concurrency | |
|---|---|---|
| What it does | Guarantees this function N units and also caps it at N | Keeps N environments initialized ahead of time; init code runs at allocation |
| Removes cold starts? | No. It only decides who gets capacity | Yes, for requests that land on those N environments |
| Cost | No extra charge | Billed per GB-second for as long as it is configured, used or not |
| Applies to | The function | A published version (an immutable, numbered snapshot of the function's code and configuration) or alias (a named pointer, like prod, that you can repoint to a different version), never $LATEST (the mutable pointer to whatever code is currently deployed, which is why it cannot be pinned to provisioned concurrency) |
| Above the setting | Throttled (429) | Overflow uses normal on-demand environments, which cold start |
| Typical use | Protect a critical function from being starved; stop a function from flooding a database | Interactive, latency-sensitive endpoints |
If both are set, provisioned concurrency cannot exceed reserved concurrency. Setting reserved concurrency to 0 is the documented way to stop a function completely.
Worked example
Assume an endpoint whose handler takes 120 ms. It normally sees 100 RPS with 12 environments warm, and the account quota is the default 1,000.
# Environments needed = RPS x duration (Little's law: this alone gives the concurrency,
# i.e. the environment count). The separate 10x-quota request-rate ceiling only binds
# for average durations under about 100 ms, which this 120 ms example never reaches.
duration_s = 0.120 # assumed average handler duration for the endpoint
bucket = 1000 # scaling rate: up to 1,000 new environments per function...
refill_per_s = 100 # ...refilled continuously (1,000 per 10 s)
warm_now = 12
def report(rps, account_limit):
need = rps * duration_s
cold = max(0, min(need, account_limit) - warm_now)
wait = max(0, cold - bucket) / refill_per_s
short = max(0, need - account_limit)
print(f"{rps:>6} RPS, quota {account_limit:>5}: need {need:>5.0f}; new (cold) environments {cold:>5.0f}; "
f"seconds beyond the first 1,000 new: {wait:>4.1f}; throttled capacity gap: {short:.0f}")
for rps in (100, 5_000, 15_000):
report(rps, 1_000) # default Regional concurrency quota
report(15_000, 3_000) # after a quota increase
pc_units, mem_gb, price_pc = 600, 1.0, 0.0000041667 # USD per GB-second, us-east-1, x86
hour = pc_units * mem_gb * 3600 * price_pc
print(f"provisioned concurrency {pc_units} x {mem_gb} GB: ${hour:.2f}/hour, ${hour * 730:,.0f}/month if always on")
100 RPS, quota 1000: need 12; new (cold) environments 0; seconds beyond the first 1,000 new: 0.0; throttled capacity gap: 0
5000 RPS, quota 1000: need 600; new (cold) environments 588; seconds beyond the first 1,000 new: 0.0; throttled capacity gap: 0
15000 RPS, quota 1000: need 1800; new (cold) environments 988; seconds beyond the first 1,000 new: 0.0; throttled capacity gap: 800
15000 RPS, quota 3000: need 1800; new (cold) environments 1788; seconds beyond the first 1,000 new: 7.9; throttled capacity gap: 0
provisioned concurrency 600 x 1.0 GB: $9.00/hour, $6,570/month if always on
What the numbers mean:
- Burst to 5,000 RPS. You need 600 environments. The scaling rate is not the problem: its token bucket starts full at 1,000 (see above), so all 588 new environments can be created immediately from that reserve, with no need to wait for the refill. The problem is that 588 requests each wait for a cold start at the same moment, and that is what p99 latency (the latency 99% of requests beat) will show. Init time depends on runtime and package size. Read it from the
Init DurationLambda logs rather than guessing. - Burst to 15,000 RPS on the default quota. You need 1,800 environments but only 1,000 can exist. About 800/1,800, or 44%, of the load is throttled.
- Same burst after raising the quota to 3,000. The first 1,000 new environments appear immediately. The remaining 788 arrive at the refill rate of about 100 per second, so over about 7.9 seconds. Requests above capacity during those seconds are still throttled.
- Cost of pre-warming 600 x 1 GB: $9.00 an hour, or about $6,570 a month if it is left on all the time. That is why you schedule it rather than leave it on.
What this means for a latency-sensitive endpoint
- Provisioned concurrency on the alias the API actually invokes (never on
$LATEST, which cannot hold provisioned concurrency). Size it to the expected burst plus a buffer. Lambda's own guidance is +10%, so 600 becomes 660. If the burst is predictable, raise it on a schedule. - Reserved concurrency at or above that number. This stops other functions in the account from starving the endpoint. Without it, a noisy batch job (an unrelated function in the same account that suddenly scales hard, for a backfill or a retry storm) in the same Region can use up the shared 1,000.
- Quota headroom requested ahead of time. Check the layers in front of Lambda too. API Gateway (AWS's managed HTTP front door) has its own default account throttle of 10,000 RPS per Region.
- Make cold starts cheaper for the overflow. Keep the package small, create SDK clients at top level so warm invocations reuse them, and defer rarely used clients. For supported runtimes, snapshot restore (Lambda SnapStart) is a cheaper alternative to keeping environments warm.
- For bursts beyond anything you pre-warmed, decide the behaviour explicitly. Either return a fast 429 and have clients retry with jittered backoff (each retry waits longer than the last, with a bit of randomness mixed in so many clients do not all retry at the same instant), or, if the work does not need an immediate answer, accept it into a queue.
Pitfalls
- Invoking
$LATESTor the wrong alias. Provisioned concurrency sits unused, you still pay for it, and cold starts continue. - Treating reserved concurrency as a performance feature. It warms nothing.
- Expecting target-tracking auto scaling (a controller that watches a metric and adds or removes capacity to hold it near a target value) to catch a spike. Application Auto Scaling on
ProvisionedConcurrencyUtilizationneeds the load to last about 3 minutes (three data points) before it adds capacity. It handles slow drift, not a sudden burst. - Forgetting the separate request-rate ceiling for sub-100-ms functions. Little's law alone gives the right environment count even for a fast function (a 20 ms function at 30,000 RPS needs only 600 environments), but request rate can still hit the platform's 10x-of-quota ceiling before concurrency does: at the default 1,000-unit quota that ceiling is 10,000 RPS, so the 30,000-RPS case above throttles despite having plenty of environment headroom.
- Scale-out has side effects downstream. 600 new environments can mean 600 new database connections. Plan the downstream protection together with the scaling.
A function is hitting downstream throttling from a database or third-party API during traffic spikes. Design a strategy to handle this gracefully: consider buffering with queues, circuit breakers, retry behavior, rate-limiting at the edge, and what a degraded-mode response looks like when a dependency is unavailable. How would you implement this on a serverless platform specifically?
Sample Answer
Direct answer
The platform's strength is the problem here. A serverless function can scale from 5 to 500 copies in seconds. Each copy has its own retry loop and its own memory, so a traffic spike turns into a flood aimed at a dependency that cannot scale. The fix is to match the rate at which you call the dependency to what it can take:
- put a queue between the trigger and the work
- cap how many function instances may call the dependency at once
- retry with exponential backoff and jitter (a small random extra delay added to each retry, so many instances retrying at once do not all fire again at the same instant) at one layer only
- trip a circuit breaker when the dependency is failing
- decide in advance what a degraded response looks like
On a serverless platform most of this is configuration: the queue trigger's maximum concurrency, reserved concurrency, the visibility timeout (the window during which a message SQS has handed to a consumer is hidden from everyone else; if it is not deleted before the window ends, it reappears on the queue) used as a backoff timer, partial batch responses, and a dead-letter queue.
1. Buffer with a queue
If the caller doesn't need the result immediately, accept the request into a queue (Amazon SQS, Simple Queue Service, here) and return 202 Accepted with a status link. The queue absorbs the spike and holds work while the dependency is down. SQS keeps messages for 4 days by default and up to 14.
2. Rate-match the consumers (the core control)
Set MaximumConcurrency on the SQS event source mapping (the piece of Lambda configuration that polls the queue on your behalf and invokes your function with each batch it reads; allowed range 2 to 1,000). This caps how many function instances the poller (that same polling component) will run at once. Set the function's reserved concurrency to at least that value. Work out the number from the dependency's limit, not from your traffic. See the worked example below.
Reserved concurrency alone is the wrong tool. The poller keeps invoking, Lambda throttles the invocations, and throttled batches go back to the queue and use up their receive count (SQS's own counter of how many times a message has been handed to a consumer without being deleted) on the way to the dead-letter queue (DLQ). MaximumConcurrency stops the poller from invoking in the first place.
If you need an exact global rate across instances (for example, "100 calls per second" rather than "20 in flight"), keep a shared token bucket (a counter that refills at a fixed rate; each call must take one token before it is allowed to proceed, which caps the combined rate no matter how many instances are calling) in a fast store such as Redis or a DynamoDB counter. The cost is one extra round trip per call.
3. Retries: one layer, with backoff and jitter
Retries can happen in five places: the HTTP/SDK client, the function code, Lambda's asynchronous retries (2 by default), SQS redelivery after the visibility timeout, and the original caller. Stacked, they multiply (see the worked example). Choose one:
- Retry only retryable failures: 429, 503 and timeouts. Honour
Retry-Afterwhen the dependency sends it. Never retry a 400. - Use the queue as the retry engine. On failure, call
ChangeMessageVisibility(the API that resets a message's visibility timeout) with a delay that grows withApproximateReceiveCount(that same receive-count attribute). This gives exponential backoff without holding a function instance open while it waits. - Enable partial batch responses (
ReportBatchItemFailures), so only the failed messages come back. - Send poison messages (messages that keep failing no matter how many times they are retried) to a DLQ (dead-letter queue) once their receive count passes
maxReceiveCount, a threshold you set to at least 5, and alarm on DLQ depth.
4. Circuit breaker
A circuit breaker counts recent failures. When they cross a threshold it moves from closed (calls flow) to open (fail immediately without calling). After a cool-down it moves to half-open (let a few probe calls through) and closes again if they succeed.
On serverless the breaker's state lives in one environment's memory, so each instance learns about the outage on its own. With concurrency capped at 20, that is acceptable: at most 20 probes per half-open window. For hundreds of instances, share the state in one small record (Redis key or DynamoDB item with a time-to-live), at the price of one read per call.
The serverless-specific trap: when the breaker is open, "fail the batch fast" puts messages back on the queue and increments their receive count. A 30-minute outage then pushes healthy messages into the DLQ. Two options:
- Pause consumption. Automation disables the event source mapping (
UpdateEventSourceMappingwithEnabled=false) when the breaker opens, and re-enables it after the cool-down. - Keep the breaker's fast failures below
maxReceiveCountby using long visibility delays.
5. Rate limiting at the edge
Reject excess traffic before it costs a function invocation:
- API Gateway stage and method throttling (rate caps set on the whole API stage, or on one method within it), plus usage plans (API Gateway's per-API-key throttle and quota configuration) with per-API-key rate and burst limits (the same token bucket mechanism as above).
- A rate-based rule in the web application firewall (AWS WAF) that limits each source IP.
- A throttled request gets a 429 and never runs your code.
6. Degraded mode when the dependency is down
| Request type | Degraded response |
|---|---|
| Read that tolerates staleness (prices, profile) | Serve the last cached value with an "as of" timestamp |
| Write the user doesn't need confirmed now (email, sync to the CRM, the customer relationship management system) | Accept, queue, return 202 and confirm later |
| Optional feature (recommendations) | Omit that section of the page |
| Critical synchronous call (payment authorization) | Fail fast with 503 and Retry-After. Don't hang until the function timeout |
Timeouts are part of the design. Give the dependency call a short client timeout (for example 2 s) that is well under the function timeout. Otherwise a hung partner keeps all 20 capped instances busy doing nothing.
Worked example
A partner API allows 100 requests per second. Each call takes 200 ms, and a function instance makes its calls one after another. Traffic arrives at 60 messages per second.
vendor_limit_rps = 100 # partner API contract
call_s = 0.200 # one partner call, sequential inside an invocation
calls_per_invocation_per_s = 1 / call_s
max_concurrency = vendor_limit_rps / calls_per_invocation_per_s
print(f"queue consumer MaximumConcurrency = {vendor_limit_rps} / {calls_per_invocation_per_s:.0f} = {max_concurrency:.0f}")
layers, retries = 3, 3
print(f"worst-case attempts reaching the partner per user request: (1+{retries})^{layers} = {(1 + retries) ** layers}")
arrival_rps, outage_s = 60, 600
backlog = arrival_rps * outage_s
drain_s = backlog / (vendor_limit_rps - arrival_rps)
print(f"backlog after a {outage_s // 60}-minute outage: {backlog:,} messages; "
f"drain time at {vendor_limit_rps}/s while {arrival_rps}/s keep arriving: {drain_s:.0f} s ({drain_s / 60:.0f} min)")
queue consumer MaximumConcurrency = 100 / 5 = 20
worst-case attempts reaching the partner per user request: (1+3)^3 = 64
backlog after a 10-minute outage: 36,000 messages; drain time at 100/s while 60/s keep arriving: 900 s (15 min)
- One instance makes 5 calls per second, so 20 instances exactly fill the 100 RPS contract. Set
MaximumConcurrency=20and reserved concurrency to at least 20. (If a handler makes its calls in parallel, divide by that fan-out as well.) - Three layers each retrying 3 times means 64 attempts reach the partner for one user request in the worst case. That is how a small blip becomes a retry storm (retries compounding across layers until the call volume reaching the dependency is far higher than the original traffic, which can itself take the dependency down).
- A 10-minute outage leaves 36,000 messages queued, well inside the retention period. Once the dependency recovers, the backlog drains in about 15 minutes, because only 40 RPS of headroom remains while new traffic keeps arriving. Add the 10-minute outage itself and the incident runs for about 25 minutes end to end, from the moment the outage starts to the moment the backlog is fully cleared, which is the number to put in front of the product owner as "how long until we are fully caught up." That is not how long any single message sits stale: SQS works through the backlog in roughly arrival order, so the very first message queued is also first in line once the dependency is back, and its own wait is close to the 10-minute outage, not 25 minutes.
flowchart LR
client[Client] --> gw[API Gateway, throttled]
gw --> q[Queue]
q --> fn[Function, max concurrency 20]
fn -->|breaker closed| partner[Partner API, 100 RPS]
fn -->|breaker open| cache[(Last known values)]
q -.->|after 5 receives| dlq[DLQ]
Pitfalls
- Scaling the function to "handle the load" when the dependency is the bottleneck. More instances produce more throttling.
- Retries at every layer, with no jitter. Every instance retries in the same second.
- In-memory rate limiters that assume a single process.
- A visibility timeout barely above the function timeout: a batch that was throttled or slow reappears while it is still being processed, producing duplicates. Follow AWS's 6x guidance: set the queue's visibility timeout to at least six times the function's timeout (plus any batching window), so a batch that gets retried after a throttle still has room to finish before it goes invisible-then-visible again. (A visibility timeout shorter than the function timeout is rejected by Lambda outright.)
- No DLQ alarm, so failed work sits unseen for days.
Your team runs a Flask-based inference server on EC2 and wants to migrate it to Lambda to cut ops overhead. Walk through the migration: how do you package the model and its native dependencies, what do you do about model size and cold starts, how do you replace whatever persistent, in-process caching the EC2 service relied on, and what benchmarking methodology and acceptance criteria would you use before cutting traffic over?
Sample Answer
Direct answer
Ship the service as a Lambda container image built on the Lambda base image for the exact architecture you will run, wrap the unchanged Flask app with an adapter, bake the model into the image and load it once during initialization, replace the in-process cache with a shared cache tier plus a small per-environment cache, and gate the cutover on a replayed-traffic benchmark with written acceptance criteria and a weighted, reversible traffic shift. The two things that surprise teams are that the in-process cache stops working the way it did, and that cold starts land on real users.
1. Packaging the model and native dependencies
- Container image, not a zip. A zip deployment is capped at 250 MB unzipped including layers (optional, pre-packaged bundles of code or libraries a function can attach without including them in its own deployment package); a CPU PyTorch or scientific stack alone approaches that. A container image can be up to 10 GB uncompressed.
- Build native libraries where they will run. Wheels (Python's pre-built package format,
.whlfiles) such asnumpy,onnxruntime, ortorchcontain compiled code,.sofiles, compiled shared-library binaries, for one CPU architecture and C library (the low-level system library, for example glibc, that compiled code links against; a mismatch between the version it was built against and the one present at runtime breaks it). Build inside the AWS Lambda Python base image for the target (x86_64orarm64), not on a laptop, so the.sofiles match. Install CPU-only builds; GPU builds pull in gigabytes of CUDA (NVIDIA's GPU libraries) you cannot use. - Keep Flask. The AWS Lambda Web Adapter runs as a Lambda extension (a companion process that runs alongside your code inside the same execution environment, commonly shipped as a layer), starts your web server inside the environment, and translates Lambda invocation events into HTTP requests, so the Flask app runs essentially unchanged. The alternative is a WSGI (Web Server Gateway Interface, Python's standard web-app interface) shim (a thin adapter that translates calls from one interface to another without touching your application code) that converts events into WSGI calls. Either way, keep the business logic independent of the event format.
- Mind the payload limits. Synchronous request and response bodies are capped at 6 MB each. If clients send images or large feature batches, move them to object storage and pass a reference.
2. Model size and cold starts
A cold start is when Lambda must create a new execution environment: download the image, start the runtime, run your initialization code, and only then run the handler.
| Where the model lives | Pros | Cons |
|---|---|---|
| Baked into the image (recommended up to a few GB) | Versioned with code; no network fetch at init | Larger image; every model change is a deploy |
Downloaded from object storage to /tmp at init | Model updates without redeploy | Download on every cold start; /tmp is up to 10,240 MB |
| Shared file system | One copy for all environments | Requires VPC attachment; adds a network dependency to init |
Load the model in global scope (module level, outside any function, so it runs once when the module is first imported and the loaded model is shared by every later call in that environment) so warm invocations reuse it. Constraints to design around:
- The INIT phase is limited to 10 seconds for on-demand functions. If it overruns, Lambda retries init during the first invocation under the function timeout, so the user pays for it. Provisioned concurrency lifts that limit and pre-initializes environments.
- INIT is billed for on-demand zip functions since August 2025 and was already billed for container images, so a slow model load costs money as well as latency.
- Memory sizing: resident model plus framework plus request working set, with headroom. CPU scales with memory (1,769 MB is one vCPU), so a 3 GB function also gets roughly 1.7 vCPUs for inference.
- Levers: shrink the model (quantize, export to ONNX, the Open Neural Network Exchange format), lazy-import rarely used modules, provisioned concurrency on the alias (a named, mutable pointer to a specific function version, so callers reference a stable name while the version behind it changes) that serves user traffic, or SnapStart (restoring new environments from a snapshot of an initialized one; supported for Python 3.12 and later) if the environment is compatible; note SnapStart does not support provisioned concurrency or
/tmpabove 512 MB.
3. Replacing the persistent, in-process cache
On EC2, one long-lived process served all requests, so a Python dict or LRU (least-recently-used) cache of feature lookups or embeddings warmed up once and stayed hot for days. On Lambda:
- Each environment handles one request at a time, so concurrency means many environments, each with its own copy of the cache.
- Environments live for minutes to hours and Lambda recycles them every few hours even under constant load, so even a function under steady, nonstop traffic sees a continuous trickle of newly cold environments, and newly empty local caches, not only at deploy time or when traffic first spikes.
Worked example: Little's law says concurrency equals arrival rate multiplied by time per request. At 40 requests per second and 150 ms per request that gives 40 × 0.15 = 6 concurrent environments; a burst to 200 requests per second needs 200 × 0.15 = 30. A key that was fetched once on EC2 may now be fetched up to 30 times before every environment has it, and each new environment starts cold.
The replacement:
- Shared second-level (L2) cache (a managed Redis or Memcached cluster, or a key-value table with TTL) that all environments read, so one miss warms it for everyone.
- Tiny first-level (L1) cache in global scope with a short TTL for the hottest keys, bounded in size so it cannot push the function out of memory.
- Precompute where possible. If the cache existed to hide an expensive feature computation, move that to a batch job writing to the key-value store.
- If the shared cache lives in a VPC, the function must be VPC-attached; account for the connection count (one client per environment, so 30 environments means about 30 connections).
4. Benchmarking methodology and acceptance criteria
Method
- Replay real traffic. Capture a representative day of requests (payload mix and arrival pattern, including idle periods and bursts) and replay against both stacks.
- Shadow mode first. Mirror production requests to Lambda without returning its answers; compare predictions request by request.
- Separate warm and cold latency. Report p50, p95 and p99 (50th, 95th and 99th percentile latency) for warm invocations, then cold-start rate and cold latency separately (the
Init Durationfield in each invocation's REPORT log line, a summary line Lambda automatically appends after every invocation, giving Duration, Billed Duration, Memory Used and, on a cold invocation, Init Duration), then the blended user-facing distribution. - Burst test from idle to peak to observe how many requests hit cold environments.
- Cost per million predictions from measured billed duration and memory, compared with the EC2 bill.
Acceptance criteria (write them before the test so nobody moves the goalposts)
| Criterion | Example gate |
|---|---|
| Prediction parity | Identical class labels for 100% of shadowed requests; numeric outputs within an agreed tolerance (different BLAS builds, the linear-algebra libraries underneath NumPy and PyTorch, change low-order bits) |
| Warm latency | Warm p99 no worse than the EC2 p99 plus an agreed margin |
| Blended latency | Blended p99 inside the service's SLO (service-level objective) including cold starts at replayed traffic |
| Errors | Error rate, and throttle rate (the fraction of requests Lambda rejects with a 429 because no concurrency was available), at or below EC2 baseline |
| Cost | Cost per million predictions at or below the target the migration was justified on |
| Downstream safety | Cache and database connection counts within their limits at burst |
Cutover: put both behind an Application Load Balancer (ALB: a managed HTTP load balancer that can route to a Lambda function as a target, the same way it routes to EC2 instances), using weighted target groups (groups of targets, here EC2 and Lambda, that each receive a configurable percentage of traffic): 5%, then 25%, 50%, 100%, holding at each step for a full traffic cycle. Rollback is a weight change. Keep the EC2 fleet for two weeks after 100%.
Pitfalls
- Building native wheels on a Mac and discovering the mismatch only in the cloud.
- Loading the model inside the handler, paying the load on every request.
- Treating a warm-only benchmark as the result, then meeting cold starts in production.
- Forgetting the cache hit rate: EC2 latency was partly cache latency, so comparing Lambda without the shared cache compares different systems.
Unlock Full Question Bank
Get access to all 27 Serverless and Function-as-a-Service Architecture interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.