AWS courseLesson 4 of 12
AWS course · Lesson 4 of 12
AWS Lambda for Data Pipelines
Use AWS Lambda in data pipelines: triggers, S3 and Kinesis processing, concurrency, cold starts, layers, memory and timeout tuning, and error handling with DLQs.
On this page
AWS Lambda runs your code in response to events without servers to manage, and bills only for the time it runs. In data pipelines it is the glue: it reacts to a file landing in S3, transforms a batch of stream records, calls an API, or starts a Glue job. It is also easy to misuse for work that is too big or too long. This lesson covers how Lambda behaves under load and failure, which is what interviews and production incidents test.
Sample code
The examples run locally: a handler is just a Python function that takes an event dict and a context object, so you can test it with hand-built events. The deployment commands are not executed here.
Lambda fundamentals
What it is. A Lambda function is your code plus configuration (runtime, memory, timeout, IAM execution role, environment variables). Each invocation runs in an execution environment, a lightweight micro-VM that Lambda creates, reuses for later invocations and eventually discards.
How it works.
- Code outside the handler (imports, client creation) runs once per environment, at initialisation. The handler runs for each invocation. Reuse clients by creating them at module level.
- An environment handles one invocation at a time. More simultaneous events mean more environments, which is what “concurrency” counts.
- The function’s execution role provides its AWS permissions, as covered in the IAM lesson.
The limits that shape data designs:
| Limit | Value |
|---|---|
| Maximum timeout | 900 seconds (15 minutes); default is 3 seconds |
| Memory | 128 MB to 10,240 MB in 1 MB steps; CPU scales with memory |
Ephemeral storage (/tmp) |
512 MB to 10,240 MB |
| Synchronous request and response payload | 6 MB each |
| Deployment package | 250 MB unzipped including layers; container images up to 10 GB |
| Default account concurrency | 1,000 concurrent executions per Region (can be raised) |
When not to use Lambda. Anything that might run longer than 15 minutes, needs more than 10 GB of memory, or needs a shuffle across many machines (large joins, big aggregations) belongs in Glue, EMR or AWS Batch. Lambda shines for small, independent units of work: one file, one message, one batch of stream records.
Pitfalls.
- Storing state in global variables or
/tmpand expecting it next time. Environments are reused unpredictably and can disappear at any moment. - Opening a database connection per invocation at high concurrency, which exhausts the database’s connection limit. Use RDS Proxy or batch writes.
In interviews. State the 15-minute limit and the memory range precisely, then explain the one-event-per-environment model, because concurrency, cold starts and costs all follow from it.
Triggers and event sources
What it is. Something has to invoke a function. Lambda has three invocation models, and each handles retries differently:
| Model | Who calls Lambda | Examples | Retries |
|---|---|---|---|
| Synchronous | The caller waits for the result | API Gateway, a direct Invoke, Step Functions task |
The caller decides |
| Asynchronous | The service hands the event to Lambda’s internal queue and moves on | S3 event notifications, SNS, EventBridge | Lambda retries (twice by default), then sends to a DLQ or on-failure destination |
| Event source mapping (polling) | Lambda polls the source and invokes with batches | SQS, Kinesis Data Streams, DynamoDB Streams, Amazon MSK, self-managed Kafka | Depends on the source: SQS redelivers; streams retry the batch until success or expiry |
How it works. For push sources (S3, SNS, EventBridge), the source needs permission to invoke the function, added to the function’s resource-based policy. For poll sources, Lambda’s event source mapping reads the source using the function’s execution role, so the role needs read permissions on the queue or stream.
import boto3
lam = boto3.client("lambda")
lam.create_event_source_mapping(
FunctionName="orders-stream-to-s3",
EventSourceArn="arn:aws:kinesis:eu-west-1:111122223333:stream/orders",
StartingPosition="LATEST",
BatchSize=500,
MaximumBatchingWindowInSeconds=5,
ParallelizationFactor=2,
BisectBatchOnFunctionError=True,
MaximumRetryAttempts=3,
MaximumRecordAgeInSeconds=3600,
FunctionResponseTypes=["ReportBatchItemFailures"],
DestinationConfig={"OnFailure": {"Destination": "arn:aws:sqs:eu-west-1:111122223333:orders-failed"}},
)
Pitfalls.
- Assuming exactly-once delivery. Every model can deliver the same event more than once, so handlers must be idempotent.
- Using EventBridge or S3 notifications where you needed buffering and back-pressure. Put SQS in between so a slow or failing consumer does not lose events.
In interviews. Name the three models and their retry behaviour. Saying “S3 invokes Lambda asynchronously, so failed events go to a DLQ or destination after two retries” shows you know where data can be lost.
Lambda with S3 and Kinesis
S3: one file at a time
What it is. The classic pattern: a file lands in raw/, S3 sends an ObjectCreated event, and a function validates or converts it and writes to curated/.
How it works. The event carries the bucket and the URL-encoded object key, not the file contents. The handler reads the object, processes it, and writes the result to a different prefix or bucket, otherwise its own output triggers it again in a loop. For bursty arrivals or when you need control over retries, route S3 events to SQS and let Lambda poll the queue.
import urllib.parse
import boto3
s3 = boto3.client("s3") # created once per execution environment
def handler(event, context):
for rec in event["Records"]:
bucket = rec["s3"]["bucket"]["name"]
key = urllib.parse.unquote_plus(rec["s3"]["object"]["key"])
if not key.startswith("raw/"):
continue # guard against recursive triggers
body = s3.get_object(Bucket=bucket, Key=key)["Body"].read()
cleaned = body.replace(b"\r\n", b"\n")
out_key = "curated/" + key.removeprefix("raw/")
s3.put_object(Bucket=bucket, Key=out_key, Body=cleaned) # deterministic key: re-runs overwrite
Kinesis: batches of stream records
What it is. An event source mapping reads each shard and invokes the function with a batch of records, in order per shard.
How it works. Record data arrives base64-encoded. If the function fails, Lambda retries the whole batch and the shard stops advancing until the batch succeeds, the record expires or the retry limits are hit. Two settings prevent one bad record (a “poison pill”) from blocking the shard: BisectBatchOnFunctionError splits a failing batch in half to isolate the bad record, and ReportBatchItemFailures lets the handler say exactly where processing failed so Lambda retries from that record only.
import base64, json
# A Kinesis event as Lambda receives it from an event source mapping (trimmed to the fields used).
def kinesis_event(payloads):
return {"Records": [
{"eventSource": "aws:kinesis",
"kinesis": {"partitionKey": p["user"], "sequenceNumber": str(1000 + i),
"data": base64.b64encode(json.dumps(p).encode()).decode()}}
for i, p in enumerate(payloads)]}
def handler(event, context=None):
"""Decode records, skip none silently, and report the first failure so Lambda retries from there."""
out = []
for rec in event["Records"]:
body = json.loads(base64.b64decode(rec["kinesis"]["data"]))
if body.get("amount") is None:
# ReportBatchItemFailures: Lambda resumes the batch from this sequence number
return {"batchItemFailures": [{"itemIdentifier": rec["kinesis"]["sequenceNumber"]}]}
out.append((body["user"], body["amount"]))
print("processed:", out)
return {"batchItemFailures": []}
good = kinesis_event([{"user": "u1", "amount": 5}, {"user": "u2", "amount": 7}])
bad = kinesis_event([{"user": "u1", "amount": 5}, {"user": "u3"}, {"user": "u2", "amount": 7}])
print(handler(good))
print(handler(bad))
processed: [('u1', 5), ('u2', 7)]
{'batchItemFailures': []}
{'batchItemFailures': [{'itemIdentifier': '1001'}]}
For a stream, report the first failed sequence number: Lambda retries from there, so the records after it are reprocessed and order is kept. Per shard, one batch is processed at a time by default; ParallelizationFactor (up to 10) processes several batches per shard concurrently while keeping order per partition key.
Pitfalls.
- Not URL-decoding S3 keys, so files with spaces or
=in their names “do not exist”. - One malformed Kinesis record blocking a shard for the whole retention period, with
IteratorAgegrowing. Set retry limits, bisecting, a maximum record age and an on-failure destination. - Writing one small S3 object per Kinesis batch, creating the small files problem. Buffer with Firehose instead when the job is “land the stream in S3”.
In interviews. Expect “a bad record is blocking your Kinesis consumer, what do you do?” Answer: partial batch responses, bisect on error, retry and age limits, and an on-failure destination that captures the bad record’s details for replay.
Concurrency and scaling
What it is. Concurrency is the number of invocations running at the same moment. It is the number that hits limits, so you need to estimate it.
How it works. By Little’s law, concurrency is roughly the request rate times the average duration:
def concurrency_needed(requests_per_second, avg_duration_seconds):
"""Little's law: in-flight executions = arrival rate x time each one takes."""
return requests_per_second * avg_duration_seconds
for rps, dur in [(50, 0.2), (200, 1.5), (1000, 3.0)]:
print(f"{rps:>5} req/s x {dur:>4}s -> about {concurrency_needed(rps, dur):>6.0f} concurrent executions")
50 req/s x 0.2s -> about 10 concurrent executions
200 req/s x 1.5s -> about 300 concurrent executions
1000 req/s x 3.0s -> about 3000 concurrent executions
The third case exceeds the default account limit of 1,000 concurrent executions per Region, so requests would be throttled unless the quota is raised. Each function can scale by up to 1,000 execution environments every 10 seconds, independently of other functions.
Controls:
- Reserved concurrency guarantees a function a slice of the account pool and also caps it. Use it to protect a downstream database (cap at what the database can handle) or to stop one function starving others. Setting it to 0 stops a function from running.
- Provisioned concurrency keeps a number of environments initialised and ready, removing cold starts for that many concurrent requests, at an extra cost.
- For SQS sources, the event source mapping’s maximum concurrency setting limits how many concurrent invocations the queue can drive.
- For Kinesis, concurrency is shards times the parallelization factor, so you scale by adding shards.
Pitfalls.
- Throttling of asynchronous invocations is retried silently for a while, so events can arrive hours late. Watch the
Throttlesmetric and the async event age. - A burst of S3 uploads (a backfill of 100,000 files) can fan out to thousands of concurrent functions that overwhelm a database. Put SQS in front and cap concurrency.
In interviews. Show the rate-times-duration estimate, quote the default 1,000 account concurrency, and explain reserved concurrency as a throttle that protects downstream systems.
Cold starts
What it is. A cold start is the extra latency when Lambda must create a new execution environment: download the code, start the runtime and run your initialisation code before the handler. Warm invocations reuse an existing environment and skip this.
How it works. Cold starts happen on the first invocation, when scaling out, and after environments are recycled. Their length depends on package size, runtime, and how much work runs at import time. For most data pipelines (asynchronous, batch-oriented) a cold start of a second or so does not matter. It matters for user-facing APIs and tight latency budgets.
Ways to reduce them:
- Keep the deployment package small; import only what you use (for example, not all of pandas for a JSON reshaping job).
- Do expensive set-up once at module level, not in the handler.
- Use provisioned concurrency for predictable low latency.
- Use SnapStart where the runtime supports it, which resumes environments from a snapshot taken after initialisation. Check the documentation for the runtimes currently supported.
- Choose arm64 (Graviton) and test whether it starts and runs faster for your code.
Pitfalls. Keeping functions warm with scheduled pings only keeps one environment warm; it does not help when traffic scales out. Use provisioned concurrency if latency truly matters.
In interviews. Define a cold start precisely (new environment plus init code), say when it matters (synchronous, latency-sensitive paths) and when it does not (most asynchronous data processing), and list provisioned concurrency and SnapStart.
Layers
What it is. A layer is a zip archive of libraries or other files that you publish once and attach to many functions. It is extracted into /opt in the execution environment, and Python layers placed under python/ are added to the import path.
How it works. A function can use up to five layers, and the function plus all layers must fit within the 250 MB unzipped deployment limit. Layers are versioned and immutable; a function references a specific layer version ARN.
mkdir -p layer/python
pip install pyarrow==21.0.0 --target layer/python \
--platform manylinux2014_x86_64 --only-binary=:all: --python-version 3.13
(cd layer && zip -r ../pyarrow-layer.zip python)
aws lambda publish-layer-version --layer-name pyarrow-21 \
--zip-file fileb://pyarrow-layer.zip --compatible-runtimes python3.13 \
--compatible-architectures x86_64
Pitfalls.
- Building native wheels on a Mac or Windows laptop and getting
ImportErrorin Lambda. Build for Lambda’s Linux platform and architecture, or in a container. - Large libraries (pandas plus pyarrow) push against the 250 MB limit. Use a container image (up to 10 GB) for heavy dependencies, or an AWS-provided layer such as the AWS SDK for pandas layer.
- Layers are not a deployment pipeline: updating a layer does not update functions until they reference the new version.
In interviews. Explain that layers share dependencies across functions and keep deployment packages small, mention the limits (five layers, 250 MB unzipped), and name container images as the answer for big dependencies.
Memory and timeout tuning
What it is. Memory is the only resource knob: Lambda allocates CPU in proportion to memory (up to 6 vCPUs at the top of the range). The timeout is the maximum time an invocation may run before Lambda stops it.
How it works. The duration part of the bill is memory multiplied by billed duration, measured in GB-seconds, plus a per-request charge. Because more memory brings more CPU, a CPU-bound function often finishes faster with more memory, and the GB-seconds can stay flat or fall. The figures below are illustrative inputs to show the calculation, not benchmarks:
def gb_seconds(memory_mb, duration_ms, invocations):
"""The duration dimension Lambda bills on: memory (GB) x billed duration (s) x invocations."""
return memory_mb / 1024 * duration_ms / 1000 * invocations
# Measured on the same workload: more memory brings more CPU, so duration drops.
trials = [(512, 4200), (1024, 2000), (2048, 1100), (4096, 1000)]
for mem, ms in trials:
print(f"{mem:>5} MB, {ms:>5} ms -> {gb_seconds(mem, ms, 1_000_000):>10,.0f} GB-s per million invocations")
512 MB, 4200 ms -> 2,100,000 GB-s per million invocations
1024 MB, 2000 ms -> 2,000,000 GB-s per million invocations
2048 MB, 1100 ms -> 2,200,000 GB-s per million invocations
4096 MB, 1000 ms -> 4,000,000 GB-s per million invocations
In this example 1,024 MB is both faster and cheaper than 512 MB, while 4,096 MB doubles the cost for almost no speed-up, because the work stopped being CPU-bound. Measure your own function at several memory sizes (the open-source AWS Lambda Power Tuning tool automates this) and read the Max Memory Used and Duration values from the REPORT line in CloudWatch Logs.
Set the timeout from observed duration with headroom, not to the 15-minute maximum “to be safe”. A short timeout fails fast and frees concurrency; a huge one lets a stuck call hold an environment, and for SQS sources the queue’s visibility timeout must be longer than the function timeout (AWS recommends at least six times) or messages reappear while still being processed.
Pitfalls.
- Picking the minimum memory to save money and getting slow, expensive runs.
- A timeout longer than the caller waits: API Gateway or a Step Functions task timing out while the function keeps running and writing.
- Running close to the memory limit, which ends in out-of-memory errors on larger inputs.
In interviews. Say that memory controls CPU too, that cost is GB-seconds plus requests, and that you tune by measuring several settings. Linking timeout to SQS visibility timeout earns extra credit.
Error handling and DLQs
What it is. Failures will happen: bad input, throttling by a downstream API, timeouts. Error handling decides whether an event is retried, parked for later or lost.
How it works by invocation model.
- Asynchronous (S3, SNS, EventBridge): Lambda retries a failed event up to two more times by default (configurable 0 to 2) and keeps events for up to six hours by default (configurable). After that it sends the event to a dead-letter queue (an SQS queue or SNS topic configured on the function) or, preferably, an on-failure destination, which also records the error and the request context.
- SQS: a failed message becomes visible again after the visibility timeout and is retried. The queue’s own redrive policy moves it to the queue’s DLQ after
maxReceiveCountattempts. ReturnbatchItemFailuresso only failed messages are retried, not the whole batch. - Kinesis and DynamoDB Streams: the batch is retried until success, the record expires, or the retry and age limits you set are reached; then details go to the on-failure destination. The data itself stays in the stream, so the destination receives metadata (shard, sequence numbers) for replay.
aws lambda put-function-event-invoke-config --function-name raw-file-validator \
--maximum-retry-attempts 1 --maximum-event-age-in-seconds 3600 \
--destination-config '{"OnFailure":{"Destination":"arn:aws:sqs:eu-west-1:111122223333:validator-failed"}}'
Make retries safe. Retries mean duplicates. Write outputs to deterministic keys, use upserts keyed by an event ID, or record processed IDs in DynamoDB with a conditional put.
Pitfalls.
- A DLQ nobody watches. Alarm on its
ApproximateNumberOfMessagesVisibleand have a documented replay procedure. - Catching every exception and returning success, which hides failures from Lambda’s retry machinery.
- Configuring a DLQ on the function for an SQS source. For SQS event source mappings, the DLQ belongs on the source queue.
In interviews. Expect “how do you make sure no events are lost?” Cover: idempotent handlers, partial batch responses, retry limits, DLQs or on-failure destinations per invocation model, alarms on them, and a replay path.
Practice questions
A nightly job processes a 40 GB file in Lambda and keeps timing out. What do you change?
Lambda’s hard limit is 15 minutes, and memory and /tmp top out at 10 GB, so this is the wrong tool. Move the job to a Glue or EMR Serverless Spark job, or to AWS Batch. If the file can be split, Lambda can process chunks in parallel through a Step Functions Distributed Map, but a single 40 GB transformation does not belong in Lambda.
What is the difference between reserved and provisioned concurrency?
Reserved concurrency sets aside part of the account’s concurrency for a function and caps the function at that number; it costs nothing extra. Provisioned concurrency keeps a number of environments initialised so those invocations have no cold start; you pay for it while it is configured. Reserved concurrency protects other systems; provisioned concurrency protects latency.
Your Kinesis-triggered function’s IteratorAge keeps growing. What could be wrong?
The consumer is falling behind: the function is too slow (tune memory, batch size, or parallelization factor), it is throttled (check concurrency limits), there are too few shards for the load, or a poison-pill record is failing and blocking a shard. Check the Errors and Throttles metrics and the logs, enable partial batch responses, bisecting and retry limits with an on-failure destination, and add shards if throughput is the issue.
An S3-triggered function writes its output to the same bucket and the bill explodes. Why?
The output write emits another ObjectCreated event that invokes the function again, which writes again: a recursive loop. Write to a different bucket or a prefix excluded by the notification filter, and guard in code by checking the key prefix. Lambda can detect and stop some recursive loops, but do not rely on it.
How do you size an SQS-triggered function’s timeout and the queue’s visibility timeout?
Set the function timeout from observed duration plus headroom. Set the queue’s visibility timeout longer than the function timeout (AWS recommends at least six times), so a message is not redelivered while still being processed. Configure a redrive policy with a DLQ after a sensible maxReceiveCount, and return batchItemFailures for partial failures.
Key takeaways
- Lambda runs one event per execution environment, up to 15 minutes, with 128 MB to 10,240 MB of memory.
- Invocation model decides retries: synchronous (caller), asynchronous (Lambda retries, then DLQ or destination), polling (source-specific).
- Concurrency is roughly rate times duration; the default account limit is 1,000 per Region, and reserved concurrency caps a function.
- Cold starts matter for latency-sensitive paths; provisioned concurrency and SnapStart address them.
- Memory also sets CPU; tune by measuring GB-seconds and duration, and keep timeouts tight.
- Make handlers idempotent and use partial batch responses, bisecting and on-failure destinations so one bad record cannot block or vanish.
Progress is saved in this browser only. No account needed.