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 an observability architecture for a serverless platform: end-to-end distributed tracing from the API gateway through the function to whatever it calls downstream, plus service-level metrics and anomaly detection, all while keeping the observability bill under control. What sampling strategy would you use, where do traces and metrics get stored, and how do you correlate telemetry across an async hop to debug one specific failed request?
Sample Answer
Direct answer
Instrument every hop with OpenTelemetry (the open standard and SDKs for traces, metrics and logs), propagate one trace context through the API gateway, the functions and every message across async hops, and put a correlation ID in every log line. Keep costs bounded with low head sampling for successful traffic plus guaranteed capture of failures, cheap aggregated metrics for SLOs (service-level objectives) and anomaly detection, and logs with short hot retention and cheap archive. To debug one failed request, walk from the request ID the client saw, to the trace ID in its logs, to the trace, across the queue via its link to the consumer's span.
Requirements I am designing to
- 100 requests per second average at the gateway (259.2M requests per month), each producing a sync function call and an async message to a worker.
- p99 (the 99th-percentile latency: the value only the slowest 1% of requests exceed) latency SLO on the API; error-budget alerting (paging when the service has spent too much of the failure rate its SLO allows, rather than on every individual error; the full mechanism is below).
- Find any specific failed request within minutes, including failures that happen after the async hop.
- Observability spend should be a small fraction of compute spend, with arithmetic below.
Architecture
flowchart LR
C[Client] --> G[API gateway]
G --> F1[API function]
F1 -->|message plus trace context| Q[Queue]
Q --> F2[Worker function]
F2 --> D[Downstream API or DB]
F1 -.spans, EMF metrics, logs.-> T[Telemetry pipeline]
F2 -.spans, EMF metrics, logs.-> T
T --> TR[Trace store]
T --> M[Metrics store]
T --> L[Log store plus archive]
Components
- Tracing: the AWS Distro for OpenTelemetry Lambda layer (AWS's own supported build of the open-source OpenTelemetry instrumentation, attached to a function without changing its code) or native AWS X-Ray (AWS's distributed tracing service) on each function. Spans (timed records of one unit of work) cover the gateway, handler, each downstream call, and the enqueue.
- Metrics: SLO metrics (requests, errors, latency histograms) emitted as structured logs in the CloudWatch embedded metric format (EMF), which the platform turns into metrics without a synchronous API call inside the function; or a Prometheus-compatible managed store (Prometheus is a widely used open-source metrics system; a managed store that speaks its query language) if the organization standardizes on it.
- Logs: structured JSON with
request_id,trace_id,function_version, outcome. Kept in hot retention (the fast, queryable store, as opposed to cheap cold archive) for 14 days, then archived to object storage. - Anomaly detection: statistical bands (a normal range computed from a metric's own recent history, so an alert fires only once it moves outside that range) on a few high-value series (error rate, p99, throughput per route), plus multi-window burn-rate alerts: an SLO like 99.9% availability allows a fixed amount of failure over its measurement window, the error budget (here, 0.1% of requests over 30 days); a burn-rate alert pages when errors are consuming that budget faster than a sustainable rate, checked over both a short window (to catch a real incident fast) and a long window (so a brief blip does not page anyone) at once, which is why it is "multi-window" (for example, page when 2% of the monthly error budget is consumed in one hour). Anomaly detection runs on aggregated metrics, never on raw traces, which keeps it cheap.
Propagating context across the async hop
The trace context (a trace ID plus parent span ID, the ID of the span that caused this one, in the W3C (World Wide Web Consortium, the standards body) traceparent header or AWS's X-Amzn-Trace-Id) must ride inside the message, because there is no HTTP connection between producer and consumer.
- SQS (Amazon Simple Queue Service) to Lambda on AWS: SQS carries the X-Ray trace header in the reserved
AWSTraceHeadermessage system attribute, and a Lambda consumer picks it up automatically; the console shows the producer's trace linked to the consumer's. - Other brokers (EventBridge, Kafka (Apache Kafka, a distributed event-streaming log), SNS fan-out, one SNS message delivered to every subscriber at once): inject
traceparentinto message attributes or headers at publish and extract it in the consumer. - Batches: one consumer invocation may process 10 messages from 10 different traces. Model the consumer's work as its own span with span links (references from one span to other related spans that are not its direct parent or children, so one span can point back to many origins) to each message's producing span, rather than pretending it has a single parent.
- Redundant safety net (belt and braces): also put a business correlation ID (order ID, request ID) in the message body and in every log line, so the chain can be reconstructed even for unsampled traces.
Sampling strategy and the cost arithmetic
Prices: X-Ray records traces at $5.00 per million, and CloudWatch Logs standard ingestion is $0.50 per GB (US East, current list).
N=100×2,592,000=259,200,000 traces per month 100% sampling=259.2×5.00=$1,2965%=$64.801%=$12.96Decision: head-sample 5% of successful requests (the decision is made at the entry point and propagated, so a trace is complete or absent, never half-recorded), and capture 100% of errors and slow requests. Capturing all errors needs tail sampling: deciding after the trace finishes. A function cannot do that alone because it sees only its own spans, so either route spans through an OpenTelemetry Collector gateway (a standalone process that receives every function's spans and applies rules to them centrally, before they reach the trace store) running tail-sampling rules (errors, latency above the SLO threshold, specific tenants), or, if you stay on head sampling only, guarantee that every error log line carries the trace ID and full context, so unsampled failures are still debuggable from logs.
Logs: 1 KB per request is 259.2 GB per month, $129.60 of ingestion. Log at INFO for outcomes only, DEBUG sampled, errors in full.
Metrics cardinality is where observability bills explode. Each unique combination of metric name and dimension values is billed as a separate custom metric ($0.30 per metric per month for the first 10,000, $0.10 for the next tier). Adding a customer_id dimension across 10,000 customers for 3 metrics creates 30,000 metrics:
That is more than all tracing at 5%. Per-customer analysis belongs in logs or traces (queried on demand), not metric dimensions.
Where telemetry is stored
| Data | Store | Retention |
|---|---|---|
| Traces | X-Ray or an OpenTelemetry-compatible trace backend | 7 to 30 days; long-term value is low |
| SLO metrics | CloudWatch or managed Prometheus | 15 months; cheap because aggregated |
| Logs | Log service, hot | 14 days, then object-storage archive queried on demand |
Debugging one specific failed request
- Client reports failure and returns the
request_idfrom the error response (always return it). - Log query on
request_idfinds the API function's line and itstrace_id. - Open the trace: gateway span, handler span, enqueue span, status OK. The failure is after the hop.
- Follow the link from the enqueue span to the worker's span (automatic for SQS to Lambda; via span links for batches).
- The worker span shows a 5 s downstream timeout; its log lines, filtered by the same
trace_id, show the payload that triggered it; the message is in the dead-letter queue, ready to redrive (resend through the same processing path for another attempt) after the fix.
If the request was not in the 5% sample and tail sampling is not deployed, steps 3 and 4 fall back to logs by trace_id and correlation ID, which is why the log discipline is non-negotiable.
Pitfalls
- Head sampling only, then discovering the failure you need was in the 95% dropped.
- Losing trace context at the queue because the producer never injected it.
- High-cardinality metric dimensions.
- Synchronous telemetry exports inside the handler that add latency and billed duration to every request.
What security best practices would you apply to a serverless function that handles sensitive data? Think about execution-role design, how secrets get to the function, network placement, and what you log, and don't log. What's different about the blast radius here compared to a traditional always-on server?
Sample Answer
Direct answer
Give each function its own narrowly scoped identity, deliver secrets at runtime from a secrets store (never baked into code or plain environment variables), put the function inside a private network only when it must reach private resources, and treat logs as a place sensitive data must never land. The blast radius (how much an attacker can reach if one component is compromised) is smaller per unit than an always-on server because each function has its own short-lived credentials and no long-lived host, but it is wider in count: a serverless app is dozens of functions, each an entry point with its own permissions to get wrong.
The four controls, one at a time
1. Execution-role design (least privilege per function)
On AWS Lambda the execution role is the IAM (AWS Identity and Access Management) role the function assumes when it runs; its temporary credentials are what your code uses to call other services.
- One role per function, not one role per app. A
get-patient-recordfunction needsdynamodb:GetItemon one table. It should not inherit thes3:PutObjectthat the export function needs. - Scope to resource ARNs (Amazon Resource Names) and actions, never
*. Add conditions where they exist (for example, restrict KMS, AWS Key Management Service, decrypt to calls made through a specific service). - Separate the invoke side too. The function's resource-based policy says who may invoke it. Lock it to the specific API gateway route (one HTTP method and path exposed by the managed service that turns HTTP requests into invocations) or queue that triggers it, so nobody else in the account can call it directly with crafted input.
2. How secrets reach the function
- Store database passwords and API keys in a secrets manager (AWS Secrets Manager, Azure Key Vault, Google Secret Manager). Grant the role read access to only that secret.
- Fetch at runtime and cache briefly. On Lambda, the AWS Parameters and Secrets Lambda Extension serves secrets from a local cache over
localhost:2773; its default cache TTL (time-to-live) is 300 seconds, adjustable withSECRETS_MANAGER_TTL. That keeps per-invocation latency and API cost down while still picking up rotations within minutes. - Why not environment variables? Lambda encrypts them at rest with KMS (AWS Key Management Service), but anyone with permission to read the function's configuration sees them in plaintext, they show up in infrastructure-as-code diffs (the reviewable, version-controlled change files that a tool like Terraform or CloudFormation generates, often posted in a pull request for anyone on the review to read), and they cannot rotate without a redeploy. Use them for non-secret configuration such as the name of the secret.
3. Network placement
- A function that only calls managed APIs (object storage, a queue, a key-value store over its public endpoint with IAM auth) does not need to be in a VPC (Virtual Private Cloud, your private network in the cloud). Putting it there adds configuration without adding protection.
- A function that talks to a private database must be attached to the VPC in private subnets (network segments with no direct route to the public internet), with a security group (a stateful firewall attached to the function's network interface, allowing only the traffic you list) that allows only the database port. Reach other AWS services through VPC endpoints (private connections to AWS APIs) rather than a NAT gateway (network address translation gateway, which gives private subnets outbound internet access), and restrict egress (outbound traffic leaving the network) so compromised code cannot exfiltrate (quietly send stolen data out to an attacker) to arbitrary internet hosts.
4. What you log, and what you never log
| Log | Never log |
|---|---|
| Request ID, trace ID, function version | Raw request bodies containing personal or payment data |
| Outcome, latency, error class | Secrets, tokens, Authorization headers |
| Hashed or tokenized user identifiers | The full event object "for debugging" |
A hashed identifier is a one-way scramble that cannot be reversed back to the original value; a tokenized identifier swaps the real value for a reference id that is meaningless without a separate, tightly controlled lookup table. Either lets you correlate one user's events across logs without the logs themselves holding personal data.
A common leak is print(event) left in from development: on an API-triggered function the event includes headers and body. Use structured logging with an allow-list of fields, and add a log-side masking policy (CloudWatch Logs data protection policies can detect and mask patterns such as card numbers) as a second net, not the first.
Worked example
A claims-processing function reads a claim record, calls a fraud-scoring API, and writes a decision.
- Role:
dynamodb:GetItemanddynamodb:UpdateItemonarn:...:table/claims,secretsmanager:GetSecretValueon the one fraud-API key secret,kms:Decrypton the table's key (DynamoDB tables can be encrypted at rest with a customer-managed KMS key instead of the AWS-owned default; when they are, reading or writing items needs decrypt or encrypt permission on that specific key too). Nothing else. - Secret: fetched via the extension at invoke time, cached up to 300 s; rotation monthly.
- Network: no VPC, because both the table and the fraud API are reached over authenticated public endpoints. If the claims store moved to a private relational database, the function would move into private subnets with a VPC endpoint for Secrets Manager.
- Logs:
claim_id_hash,decision,fraud_score_bucket,latency_ms. The claimant's name and bank details never appear.
Blast radius: serverless versus an always-on server
| Dimension | Always-on server | Serverless function |
|---|---|---|
| Credentials | One host role shared by every app and cron job on the box | One role per function, temporary credentials expiring on their own |
| Persistence for an attacker | Can install a backdoor and wait | Environment is recycled; no host to persist on |
| What one exploit reaches | Everything the host role can reach, plus local files and neighbours | Only what that function's role allows |
| Number of entry points | A few | Many: every trigger (HTTP route, queue, bucket event) is one |
The honest caveats that make this a senior answer:
- Credentials are still stealable while the environment is alive. The execution role's temporary credentials are exposed to the runtime as environment variables, which means any code running inside that process, including a compromised dependency, can read them simply by reading its own process environment; nothing stops it. A dependency with remote code execution (a bug that lets an attacker run arbitrary code inside your process, for example through unsafe deserialization) or a server-side request forgery bug (a flaw that tricks the server into making an HTTP request of the attacker's choosing, for example to an internal endpoint that hands back these same credentials) can read and use them until they expire. Least privilege is what bounds the damage, not the ephemerality.
- Warm environments are reused. Globals and files in
/tmpsurvive between invocations on the same environment. Never cache one caller's sensitive data in a global or/tmpwhere the next request could read it. - Event injection is the new input surface. A function triggered by a file upload or queue message must validate that payload as strictly as an HTTP body.
Pitfalls
- One shared "lambda-role" with
AdministratorAccessbecause it was faster during a prototype. - Secrets in environment variables "because they are encrypted".
- Putting every function in a VPC by default, then opening a NAT gateway to the whole internet so they can reach public APIs.
- Logging the full event, which quietly turns the log store into the largest sensitive-data store in the account.
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.
You're deploying a function that reads an object from cloud storage, looks up data in a managed key-value store, and pulls a credential from a secrets service. Design a least-privilege IAM policy for it: how would you structure the roles and permission boundaries? Separately, how would you handle secret rotation and caching so you're not paying a cold-start latency penalty on every credential fetch, without compromising security?
Sample Answer
Direct answer
Give the function its own execution role (the IAM role it runs as, where IAM is AWS Identity and Access Management). That role allows exactly four things:
- read objects under one Amazon S3 (object storage) bucket prefix
- read items from one table
- read one secret
- decrypt with that secret's KMS (AWS Key Management Service) key, and only when the request comes through Secrets Manager. Secrets Manager encrypts every secret at rest with a KMS key, so even though the code only calls
secretsmanager:GetSecretValue, Secrets Manager in turn needs the caller allowed to decrypt with that key, or the call fails
Scope each permission to a specific resource ARN (Amazon Resource Name, the unique id of an AWS resource), not *. Wrap every role your deployment pipeline creates in a permissions boundary: a ceiling policy attached to the role that caps what it can ever be granted, so even a mistaken policy attached later still cannot grant more than the ceiling allows.
For secrets, fetch once per execution environment (the sandboxed instance of your function that a cold start creates, and that later requests reuse while it stays warm) and cache in memory with a short TTL (time-to-live), either through the Lambda secrets extension or a library cache. That way warm invocations (requests that land on an already-running environment instead of triggering a new cold start) never pay for a fetch. Make rotation safe for cached copies by using the alternating-users rotation strategy plus refresh-on-authentication-failure.
Structuring the roles
An identity policy is attached to a principal, here the function's role, and says what that principal may call. A resource-based policy is attached to the resource itself, here the function, the secret, or the KMS key, and says who may reach that resource. A call is only allowed when both sides agree.
| Piece | Controls | For this function |
|---|---|---|
| Execution role (identity policy) | What the function may call | The four statements below, plus writing to its own log group |
| Function resource-based policy | Who may invoke the function | Only the specific bucket notification or API that triggers it |
| Permissions boundary | The maximum the execution role can ever have | Allows only S3, DynamoDB, Secrets Manager, KMS and CloudWatch Logs actions on this app's resources. Denies all iam:* |
| Deployment role (the CI, continuous integration, pipeline) | Who may create or alter roles | May create roles only if the boundary is attached, enforced with a condition key (an extra clause on a permission, here iam:PermissionsBoundary, that must also hold for the action to be allowed), and may iam:PassRole only for roles under this app's path. Passing a role means handing a role to an AWS service so it can act as that role; a deployer needs it to attach the execution role to the function it creates, and it is restricted so the deployer cannot hand out a more powerful role than this app is meant to have |
| Resource policies on the other side | Defense in depth (stacking independent controls so that if one fails or is misconfigured, another still stops the request) | The secret's resource policy and the KMS key policy name this role; the bucket policy restricts the prefix |
Effective permission is the intersection of the identity policy, the boundary, and any organization-level service control policies (guardrails set above the account, by the AWS Organization, that cap what any role in the account can ever do no matter what its own policy allows). An explicit deny anywhere wins.
(Each statement below has a Sid, a statement id: just a human-readable label. AWS ignores it when deciding access; only Effect, Action, Resource and Condition do.)
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "ReadInputObjects",
"Effect": "Allow",
"Action": "s3:GetObject",
"Resource": "arn:aws:s3:::orders-input/incoming/*"
},
{
"Sid": "ReadCustomerLookupTable",
"Effect": "Allow",
"Action": ["dynamodb:GetItem", "dynamodb:BatchGetItem"],
"Resource": "arn:aws:dynamodb:us-east-1:111122223333:table/customer-lookup"
},
{
"Sid": "ReadOnePartnerCredential",
"Effect": "Allow",
"Action": "secretsmanager:GetSecretValue",
"Resource": "arn:aws:secretsmanager:us-east-1:111122223333:secret:prod/order-enricher/partner-api-??????"
},
{
"Sid": "DecryptOnlyViaSecretsManager",
"Effect": "Allow",
"Action": "kms:Decrypt",
"Resource": "arn:aws:kms:us-east-1:111122223333:key/1234abcd-12ab-34cd-56ef-1234567890ab",
"Condition": {
"StringEquals": { "kms:ViaService": "secretsmanager.us-east-1.amazonaws.com" }
}
}
]
}
Details that separate a real least-privilege policy from a plausible one:
- Read-only verbs. A verb is IAM's term for the action itself, the part of the permission that says what you may do, as opposed to the resource it applies to. No
PutObject,PutItemorScan.GetItemandBatchGetItemcover key lookups. AddQueryonly if the code does range reads. - Secret ARN ends in
-??????. Secrets Manager appends six random characters to secret ARNs.??????matches exactly those six and nothing longer, while a*would also matchpartner-api-admin-.... kms:ViaServicemeans the function can use the key only through Secrets Manager, not decrypt arbitrary data with it. If the bucket uses SSE-KMS (server-side encryption with a KMS key), add a second decrypt statement scoped to S3.- Logs: the basic execution permissions, scoped to this function's log group. A function inside a VPC (Virtual Private Cloud, a private network) also needs the network-interface permissions.
- Multi-tenant tables: add the
dynamodb:LeadingKeyscondition (a restriction on the table's partition key, the attribute DynamoDB uses to decide which physical partition an item lives on and the first part of every lookup) so the role can only read partition keys belonging to the caller's tenant. - Checking for unused grants: generate a policy from recorded CloudTrail (AWS's API audit log) activity with IAM Access Analyzer. Its unused-access findings flag permissions the function has but never uses.
Secret caching without the cold-start penalty
Where the cost comes from. Calling Secrets Manager inside the handler adds a network round trip to every invocation and an API charge per call. The fix is to fetch once per execution environment and reuse the value across warm invocations.
Two supported ways to do it:
- The AWS Parameters and Secrets Lambda Extension, added as a layer (a packaged bundle of code or binaries a function attaches without bundling it into its own deployment package). The function calls
http://localhost:2773/secretsmanager/get?secretId=...with the headerX-Aws-Parameters-Secrets-Tokenset to the function's session token (the temporary credential Lambda issues alongside the execution role's access key, proving to the local extension that the caller really is this running function). It caches for 300 s by default (configurable 0 to 300 s throughSECRETS_MANAGER_TTL), holding up to 1,000 secrets. - Powertools for AWS Lambda parameters with a
max_age. Or a small module-level cache of{value, fetched_at}.
When the fetch happens. Pre-warmed environments (provisioned concurrency) run init code before any request, so fetching at init moves the cost off the request path entirely. With snapshot restore (SnapStart), anything fetched during init is baked into the snapshot and shared by every restored environment. So fetch after restore, not during init.
Never put the secret in an environment variable (readable by anyone with lambda:GetFunctionConfiguration), in /tmp, or in logs. Note that AWS's own sample code prints the retrieved secret. Remove that line.
Cost:
price_per_call = 0.05 / 10_000 # Secrets Manager API calls, USD
invocations_per_s, envs, ttl_s = 200, 50, 300
month_s = 30 * 24 * 3600
no_cache = invocations_per_s * month_s
cached = envs * (month_s / ttl_s) # worst case: every environment refreshes every TTL
print(f"fetch per invocation: {no_cache:,.0f} calls/month = ${no_cache * price_per_call:,.2f}")
print(f"cache per environment, {ttl_s} s TTL, {envs} environments: {cached:,.0f} calls/month = ${cached * price_per_call:,.2f}")
fetch per invocation: 518,400,000 calls/month = $2,592.00
cache per environment, 300 s TTL, 50 environments: 432,000 calls/month = $2.16
At 200 invocations per second, fetching per invocation costs about $2,592 a month in API calls alone (plus the latency). Caching per environment with a 300 s TTL, even assuming all 50 environments refresh every 5 minutes, costs $2.16.
Rotation that is safe for cached values
Secrets Manager rotation moves three staging labels between versions of the secret: AWSCURRENT marks the version callers get by default, AWSPENDING marks the new version created for an in-progress rotation and not yet promoted, and AWSPREVIOUS marks the version that was AWSCURRENT immediately before the last rotation.
- Single-user rotation changes the password of the one database user in place. Every environment holding the old value in its cache then fails authentication for up to one TTL (up to 5 minutes).
- Alternating-users rotation keeps two users. Rotation changes the password of the inactive user and then makes it
AWSCURRENT. The previous credential stays valid until the next rotation, so cached copies keep working. Use this. - Refresh on authentication failure: if a call fails with an auth error, bypass the cache (for example Powertools'
force_fetch=True, or a direct SDK call), fetchAWSCURRENTagain, and retry once. - TTL choice: with alternating users, the TTL no longer has to be short for rotation's sake. It is bounded by how quickly you must stop using a compromised secret. For an emergency revoke, rotate and then publish a new function version or move an alias (a named, mutable pointer to one specific version, used so callers reference a stable name while the version behind it changes) to recycle the environments: pointing callers at a new version or alias forces Lambda to start fresh execution environments for it rather than reuse existing warm ones, so any environment still holding the old cached secret in memory simply stops receiving traffic, rather than waiting out the cache TTL.
Pitfalls
- One shared role for "all the app's functions". One compromised function gets every permission.
secretsmanager:GetSecretValueon*, orkms:Decryptwith no conditions.- Caching forever with no TTL, combined with single-user rotation: an outage every rotation.
- Fetching during snapshot init, or logging the secret "just for debugging".
A long-running job (say, a Kubernetes job, or a managed batch or training service) needs to be triggered, tracked, and retried from your serverless layer. How do you pass parameters in, track progress, handle retries and idempotency if it fires twice, and notify on completion? Where does the workflow's own state actually live?
Sample Answer
Direct answer
The serverless function should start the job and then get out of the way. It should never wait for the job. A durable workflow engine (AWS Step Functions Standard, or equivalents such as Azure Durable Functions or Google Workflows) owns the run. The workflow's execution holds the state: the input parameters, the current step, retry counts and results are all stored in its execution history, not in any function's memory.
- A starter (a function or an event rule) starts an execution named after the request id, which makes a duplicate trigger harmless.
- The workflow submits the job through a "run a job and wait" integration, or a callback token (a one-time credential the workflow hands an external worker so that worker, not the workflow, decides when the step is done; covered in detail below).
- It retries defined failures with backoff (each retry waits longer than the last, instead of hammering the dependency again immediately).
- It publishes a completion or failure event that notifies whoever needs to know.
Why not keep this in a function
- A Lambda invocation lasts at most 15 minutes. Training and batch jobs run for hours.
- A function that polls for hours is paying to wait.
- If the environment is recycled, whatever it held in memory is gone.
A Standard workflow can run for up to a year, and it records every state transition.
Passing parameters in
- The execution input (JSON, up to 256 KiB: a kibibyte is 1,024 bytes, so about 256,000 bytes) carries the request id, the job type, and pointers to the data: input and output URIs.
- Large configuration or datasets travel by claim check: store them in object storage and pass the URI.
- The workflow maps the input into the job's own parameters (container environment variables, a Kubernetes Job spec, training hyperparameters).
Tracking progress
| Signal | Mechanism | Granularity |
|---|---|---|
| Which step the run is on | The workflow's execution status and history | Coarse, always available |
| Did the job finish | A .sync integration (the .sync suffix appended to the Resource ARN, shown below, tells Step Functions to submit the job and then wait, instead of firing it and moving straight to the next step): Step Functions waits and moves on when the job reaches a final state | Coarse |
| Percent complete, current epoch | The job writes a progress record (a key-value table item keyed by request id) that the UI reads | Fine |
| Is the job still alive | Callback pattern with a task token (a unique string Step Functions generates for this step and hands to whoever will eventually finish it; the workflow pauses right there until that same token is handed back): the job calls SendTaskHeartbeat (an API meaning "still working, do not time me out") on a schedule, and HeartbeatSeconds (the longest gap Step Functions allows between heartbeats) fails the task if heartbeats stop | Liveness |
Integration details matter.
- Step Functions' Amazon EKS (Elastic Kubernetes Service, managed Kubernetes)
runJob.syncintegration detects completion by polling: about once a minute at first, slowing to about once every 5 minutes. So completion can be noticed several minutes late. - The EKS integration does not support the callback pattern (the task-token approach above), and it only works with clusters whose Kubernetes API endpoint is public.
- For a private cluster, use callback instead. The workflow sends a message containing its task token to a queue:
sqs:sendMessage.waitForTaskTokenis the Resource suffix that tells Step Functions "send this message, then pause the state machine right here until someone calls back with this token." An in-cluster worker (something running inside the private cluster, able to reach it) picks the message off the queue, runs the Kubernetes job, and callsSendTaskSuccessorSendTaskFailure, passing the same token back, to tell Step Functions which paused execution to resume and how. - AWS Batch (a managed batch-job scheduler) and Amazon SageMaker (a managed machine learning platform) training jobs support
.syncdirectly.
Retries
- Infrastructure failures (spot capacity reclaimed: the cloud provider taking back the cheaper, interruptible compute it lent you, or a node lost) are retried closest to the job: AWS Batch's job
retryStrategy(how many times Batch itself resubmits the same job attempt before giving up), or KubernetesbackoffLimit(the same idea for a Kubernetes Job: how many failed pod attempts it tolerates). - Workflow-level
Retrycovers failures of the task itself, withIntervalSeconds,BackoffRateandMaxAttempts. - Bad input is not retried. A
Catchroutes it to the failure notification. - A
TimeoutSecondson the job state turns a job that hangs silently into a visible failure.
This is written in the Amazon States Language, the JSON-based language Step Functions state machines are defined in. A few conventions used below: a field ending in .$ (like JobName.$) takes a JSONPath expression instead of a literal value, so "$.requestId" means "read the requestId field out of this state's input," not the literal string $.requestId. Resource names the action to call using an ARN (Amazon Resource Name: AWS's standard way of pointing at one specific resource or action); the .sync on the end of arn:aws:states:::batch:submitJob.sync is what makes this state submit the job and wait, as described above. ResultPath: "$.job" tells Step Functions to write this state's output into a new job field alongside the existing input, instead of replacing the input outright, so later states can still read $.requestId and everything else. In plain terms, the numbers below mean: give the job up to TimeoutSeconds: 43200 (12 hours) to finish; on a States.TaskFailed error, wait IntervalSeconds: 60 (1 minute) and retry, and if it fails again wait 60 x BackoffRate: 2.0 = 120 seconds before the second and last retry (MaxAttempts: 2); anything else (States.ALL) skips retrying and goes straight to NotifyFailure. The two NotifySuccess/NotifyFailure states each publish to an SNS topic (Amazon Simple Notification Service, a managed pub/sub notification service; TopicArn is that topic's ARN).
{
"Comment": "Run one job per request and notify",
"StartAt": "RunJob",
"States": {
"RunJob": {
"Type": "Task",
"Resource": "arn:aws:states:::batch:submitJob.sync",
"Parameters": {
"JobName.$": "$.requestId",
"JobQueue": "arn:aws:batch:us-east-1:111122223333:job-queue/training",
"JobDefinition": "arn:aws:batch:us-east-1:111122223333:job-definition/train:3",
"ContainerOverrides": {
"Environment": [
{ "Name": "INPUT_URI", "Value.$": "$.inputUri" },
{ "Name": "OUTPUT_URI", "Value.$": "$.outputUri" }
]
}
},
"ResultPath": "$.job",
"TimeoutSeconds": 43200,
"Retry": [
{ "ErrorEquals": ["States.TaskFailed"], "IntervalSeconds": 60, "BackoffRate": 2.0, "MaxAttempts": 2 }
],
"Catch": [
{ "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "NotifyFailure" }
],
"Next": "NotifySuccess"
},
"NotifySuccess": {
"Type": "Task",
"Resource": "arn:aws:states:::sns:publish",
"Parameters": {
"TopicArn": "arn:aws:sns:us-east-1:111122223333:job-events",
"Message.$": "States.Format('request {} succeeded', $.requestId)"
},
"End": true
},
"NotifyFailure": {
"Type": "Task",
"Resource": "arn:aws:states:::sns:publish",
"Parameters": {
"TopicArn": "arn:aws:sns:us-east-1:111122223333:job-events",
"Message.$": "States.Format('request {} failed', $.requestId)"
},
"Next": "Failed"
},
"Failed": { "Type": "Fail" }
}
}
Idempotency when the trigger fires twice
Handle it at three levels:
- The execution name is the request id. For Standard workflows, starting an execution with a name and input that match a running execution returns that same execution. If the execution has already closed, the call fails with
ExecutionAlreadyExists, which the starter treats as "already done". Express workflows (a cheaper, higher-throughput Step Functions mode built for short-lived, high-volume executions, with at-least-once semantics and no long execution history) do not give this guarantee, which is one more reason to use Standard here. - The job name or output path is deterministic. Reusing the request id as the job name (and as the output prefix) means a manual re-run overwrites rather than duplicates. Some services enforce this for you: a SageMaker training job name must be unique in the account and Region, so a second create with the same name is refused.
- Consumers of the completion event can deduplicate on the request id, because the notification itself can also arrive twice.
Notifying on completion
- The final states publish to a topic or event bus, as in the definition above.
- Step Functions also emits execution status-change events to EventBridge, so downstream systems can subscribe without the workflow knowing about them.
- Notify on failure as well as success. A job that fails silently is the common incident.
Where the state actually lives
| State | Lives in |
|---|---|
| Request record (who asked, for what) | A database table keyed by request id, for API lookups |
| Workflow progress, retries, step results | The workflow execution's history |
| Job-internal progress and checkpoints | The job's own storage (object storage checkpoints) plus the progress record |
| Logs | The job's log stream |
| In any function | Nothing that must survive the invocation |
sequenceDiagram
participant API as Starter function
participant WF as Workflow execution
participant JOB as Batch or Kubernetes job
participant BUS as Topic or event bus
API->>WF: StartExecution(name = requestId, input)
WF->>JOB: submit job and wait
JOB-->>WF: final state SUCCEEDED or FAILED
WF->>BUS: publish completion or failure
Pitfalls
- A function polling job status in a loop. It costs money, times out at 15 minutes, and loses its state.
- Random execution names. Every duplicate trigger then starts a new, expensive job.
- Retrying everything, including jobs that failed on bad input.
- Stopping the workflow and assuming the job stopped. Step Functions only makes a best-effort attempt to cancel
.syncjobs, so a stopped execution can leave a job running and billing.
Unlock Full Question Bank
Get access to all 15 Serverless and Function-as-a-Service Architecture interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.