AWS courseLesson 9 of 12
AWS course · Lesson 9 of 12
AWS Step Functions for Data Pipeline Orchestration
Orchestrate AWS data pipelines with Step Functions: state machines, Standard vs Express, Map and Parallel, retries and catches, Glue and Lambda integrations.
On this page
AWS Step Functions runs workflows defined as state machines: a JSON document that says which step runs next, what to do on errors, and how to fan work out in parallel. For AWS data pipelines it is the serverless alternative to Airflow: it starts Glue jobs, Lambda functions, Athena queries and EMR steps, waits for them, retries them and alerts on failure, with no scheduler to run. This lesson builds a real pipeline definition and runs it through a small local interpreter so you can see the retry and catch logic work.
Sample pipeline
The running example is a daily orders pipeline: list new files, validate each one in parallel, then run a Glue job and start a crawler side by side, alerting through SNS if anything fails. The definition is written in Amazon States Language (ASL). It is stored here as a Python string so the local simulation further down can load it; the JSON itself is exactly what you would deploy.
ASL = r"""
{
"Comment": "Daily orders pipeline: validate new files, run Glue, refresh the catalogue, alert on failure",
"StartAt": "ListNewFiles",
"States": {
"ListNewFiles": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": {
"FunctionName": "arn:aws:lambda:eu-west-1:111122223333:function:list-new-files",
"Payload": { "prefix.$": "$.prefix" }
},
"ResultSelector": { "files.$": "$.Payload.files" },
"ResultPath": "$.listing",
"Next": "ValidateFiles"
},
"ValidateFiles": {
"Type": "Map",
"ItemsPath": "$.listing.files",
"MaxConcurrency": 5,
"ItemProcessor": {
"ProcessorConfig": { "Mode": "INLINE" },
"StartAt": "ValidateFile",
"States": {
"ValidateFile": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": {
"FunctionName": "arn:aws:lambda:eu-west-1:111122223333:function:validate-file",
"Payload": { "key.$": "$" }
},
"Retry": [
{
"ErrorEquals": ["Lambda.TooManyRequestsException", "Lambda.ServiceException"],
"IntervalSeconds": 2,
"MaxAttempts": 3,
"BackoffRate": 2.0,
"JitterStrategy": "FULL"
}
],
"End": true
}
}
},
"ResultPath": null,
"Next": "TransformAndCatalogue"
},
"TransformAndCatalogue": {
"Type": "Parallel",
"Branches": [
{
"StartAt": "RunGlueJob",
"States": {
"RunGlueJob": {
"Type": "Task",
"Resource": "arn:aws:states:::glue:startJobRun.sync",
"Parameters": {
"JobName": "raw-to-curated-orders",
"Arguments": { "--run_date.$": "$.run_date" }
},
"Retry": [
{
"ErrorEquals": ["Glue.ConcurrentRunsExceededException"],
"IntervalSeconds": 30,
"MaxAttempts": 5,
"BackoffRate": 2.0,
"MaxDelaySeconds": 300
}
],
"End": true
}
}
},
{
"StartAt": "StartCrawler",
"States": {
"StartCrawler": {
"Type": "Task",
"Resource": "arn:aws:states:::aws-sdk:glue:startCrawler",
"Parameters": { "Name": "raw-orders-crawler" },
"End": true
}
}
}
],
"ResultPath": "$.results",
"Catch": [
{ "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "NotifyFailure" }
],
"Next": "Done"
},
"NotifyFailure": {
"Type": "Task",
"Resource": "arn:aws:states:::sns:publish",
"Parameters": {
"TopicArn": "arn:aws:sns:eu-west-1:111122223333:pipeline-alerts",
"Message.$": "$.error.Cause"
},
"Next": "PipelineFailed"
},
"PipelineFailed": { "Type": "Fail", "Error": "PipelineFailed", "Cause": "See the alert for details" },
"Done": { "Type": "Succeed" }
}
}
"""
State machines
What it is. A state machine is a set of named states and the transitions between them. Each execution receives a JSON input, passes JSON from state to state, and ends in success or failure. Every execution is recorded, so you can see which state failed, with what input and error.
How it works. The main state types:
| State | Does |
|---|---|
Task |
Calls a service: Lambda, Glue, Athena, EMR, SNS, DynamoDB, or any AWS SDK API |
Choice |
Branches on conditions in the data |
Wait |
Pauses for a time or until a timestamp |
Parallel |
Runs fixed branches at the same time and waits for all |
Map |
Runs the same steps for each item in a list |
Pass |
Passes or reshapes data |
Succeed, Fail |
End the execution |
Data flows through JSONPath fields: InputPath and Parameters shape a state’s input (keys ending in .$ are paths into the data), ResultSelector trims the service response, ResultPath decides where the result is placed in the state’s input (null discards it), and OutputPath filters what moves on. Newer definitions can set "QueryLanguage": "JSONata" and use JSONata expressions with Arguments and Output instead; the examples here use the default JSONPath.
aws stepfunctions create-state-machine --name orders-daily \
--definition file://orders-daily.asl.json \
--role-arn arn:aws:iam::111122223333:role/sfn-orders \
--type STANDARD
aws stepfunctions start-execution \
--state-machine-arn arn:aws:states:eu-west-1:111122223333:stateMachine:orders-daily \
--name orders-2026-10-05 \
--input '{"prefix": "raw/orders/", "run_date": "2026-10-05"}'
Pitfalls.
- State input and output are limited to 256 KiB. Pass S3 keys and IDs between states, not datasets.
- Forgetting
ResultPath, so a task’s response replaces the whole state and later states loserun_date. - An execution role without permission for each service the definition calls (including
glue:GetJobRunfor.syncpolling andevents:permissions for managed rules).
In interviews. Describe the state types, the JSON data flow (Parameters, ResultSelector, ResultPath) and the 256 KiB payload limit, and say that every execution’s history is visible for debugging.
Standard vs Express
What it is. Step Functions has two workflow types with very different limits and semantics. You choose when creating the state machine and cannot change it later.
| Standard | Express | |
|---|---|---|
| Maximum duration | 1 year | 5 minutes |
| Execution semantics | Exactly once per state (unless you add Retry) | Asynchronous: at least once; synchronous: at most once |
| Execution history | Stored by the service, up to 25,000 events, viewable and queryable | Sent to CloudWatch Logs if enabled |
.sync and .waitForTaskToken integrations |
Supported | Not supported |
| Pricing dimension | State transitions | Number of requests plus duration and memory |
| Typical use | Batch pipelines, long jobs, human approval, non-idempotent steps | High-volume, short event processing, streaming transforms, micro-batches |
How to choose. Data pipelines that wait for Glue or EMR jobs are Standard workflows: they run for minutes to hours, need .sync, and benefit from exactly-once steps and full history. Express suits work like “for each event from EventBridge, validate it and write to DynamoDB”, which must finish quickly, runs at very high rates, and is idempotent. Distributed Map can run its child workflows as Express executions to get cheap, massive fan-out inside a Standard parent.
Pitfalls.
- Choosing Express for a pipeline that calls Glue: it cannot use
.sync, and five minutes is too short. - Long-running Standard workflows with loops that exceed 25,000 history events. Use a Distributed Map or start a new execution to continue.
In interviews. Quote the two durations (one year and five minutes), the semantics (exactly-once versus at-least-once or at-most-once), and the pricing dimensions, then pick Standard for batch orchestration.
Map and Parallel states
What it is. Parallel runs a fixed set of different branches concurrently. Map runs the same sub-workflow for each item of a list. Both wait for all work to finish, and both fail if any branch or iteration fails (unless caught).
How it works.
- Inline Map runs iterations inside the parent execution.
MaxConcurrencylimits parallelism (0 means as many as possible, up to the inline limit of 40 concurrent iterations), and all iterations count toward the parent’s 256 KiB payload and 25,000-event history. - Distributed Map (
"ProcessorConfig": {"Mode": "DISTRIBUTED"}) launches each iteration, or each batch of items, as a separate child execution, up to 10,000 in parallel. It can read items directly from S3 with anItemReader(a list of objects, a CSV or JSON file, or an S3 Inventory manifest), group items withItemBatcher, write results to S3 withResultWriter, and tolerate a percentage of failures.
{
"Type": "Map",
"ItemReader": {
"Resource": "arn:aws:states:::s3:listObjectsV2",
"Parameters": { "Bucket": "example-lake", "Prefix": "raw/orders/dt=2026-10-05/" }
},
"ItemBatcher": { "MaxItemsPerBatch": 100 },
"MaxConcurrency": 500,
"ToleratedFailurePercentage": 1,
"ItemProcessor": {
"ProcessorConfig": { "Mode": "DISTRIBUTED", "ExecutionType": "EXPRESS" },
"StartAt": "ConvertBatch",
"States": {
"ConvertBatch": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "convert-to-parquet", "Payload.$": "$" },
"End": true
}
}
},
"ResultWriter": {
"Resource": "arn:aws:states:::s3:putObject",
"Parameters": { "Bucket": "example-sfn-results", "Prefix": "orders-convert/" }
},
"End": true
}
Pitfalls.
- Unbounded concurrency hammering a downstream system (a database, an API, Lambda’s account concurrency). Set
MaxConcurrency. - Inline Map over thousands of items: history and payload limits break it. Switch to Distributed mode.
- A single failing item failing the whole Map. Use
ToleratedFailurePercentageor a Catch inside the iteration that records the failure and succeeds.
In interviews. Explain fan-out and fan-in: Parallel for different tasks, Map for the same task over many items, and Distributed Map with S3 item readers for very large fan-outs such as processing every file in a prefix.
Error handling and retries
What it is. Any Task, Parallel or Map state can have Retry rules (try again with backoff) and Catch rules (move to a recovery state when retries are exhausted). Errors have names: service errors such as Lambda.TooManyRequestsException or Glue.ConcurrentRunsExceededException, and built-in ones such as States.Timeout, States.TaskFailed and the wildcard States.ALL.
How it works. A retrier has ErrorEquals, IntervalSeconds (default 1), MaxAttempts (default 3), BackoffRate (default 2.0), and optionally MaxDelaySeconds and JitterStrategy (FULL randomises each wait so many executions do not retry in lockstep). The wait before retry n is the interval times the backoff rate to the power n minus one, capped by the maximum delay:
def retry_delays(interval_seconds=2, backoff_rate=2.0, max_attempts=4, max_delay_seconds=None):
"""Wait before each retry in a Step Functions Retrier (JitterStrategy NONE)."""
delays = []
for attempt in range(max_attempts):
d = interval_seconds * backoff_rate ** attempt
if max_delay_seconds is not None:
d = min(d, max_delay_seconds)
delays.append(d)
return delays
print("defaults (1s, x2, 3 attempts):", retry_delays(1, 2.0, 3))
print("Glue throttling retrier :", retry_delays(30, 2.0, 5, max_delay_seconds=300))
defaults (1s, x2, 3 attempts): [1.0, 2.0, 4.0]
Glue throttling retrier : [30.0, 60.0, 120.0, 240.0, 300]
The simulation below interprets the sample definition with mocked tasks. In the first run the Glue job is “busy” twice and the retrier recovers; in the second it stays busy, the retries run out, the Parallel state’s Catch routes to the SNS alert, and the execution ends in the Fail state:
import copy
class TaskFailed(Exception):
def __init__(self, error, cause=""):
super().__init__(error)
self.error, self.cause = error, cause
def get_path(data, path):
if path == "$":
return data
for part in path[2:].split("."):
data = data[part]
return data
def set_path(data, path, value):
if path is None:
return data # ResultPath null: discard the result, keep the input
if path == "$":
return value
data = copy.deepcopy(data)
target, parts = data, path[2:].split(".")
for part in parts[:-1]:
target = target.setdefault(part, {})
target[parts[-1]] = value
return data
def render(template, data):
"""Apply Parameters / ResultSelector: keys ending in .$ are JSONPath lookups."""
if isinstance(template, dict):
out = {}
for k, v in template.items():
if k.endswith(".$"):
out[k[:-2]] = get_path(data, v)
else:
out[k] = render(v, data)
return out
return template
def matches(error, error_equals):
return "States.ALL" in error_equals or error in error_equals or (
"States.TaskFailed" in error_equals and error != "States.Timeout")
def run_task(state, data, tasks, log, name):
params = render(state.get("Parameters", {}), data)
retriers = state.get("Retry", [])
attempts = [0] * len(retriers)
while True:
try:
return tasks[state["Resource"]](params)
except TaskFailed as e:
for i, r in enumerate(retriers):
if matches(e.error, r["ErrorEquals"]):
if attempts[i] < r.get("MaxAttempts", 3):
attempts[i] += 1
log.append(f"{name}: {e.error}, retry {attempts[i]}")
break
raise
else:
raise
def run(machine, data, tasks, log):
name = machine["StartAt"]
while True:
state = machine["States"][name]
kind = state["Type"]
try:
if kind == "Task":
result = run_task(state, data, tasks, log, name)
if "ResultSelector" in state:
result = render(state["ResultSelector"], result)
data = set_path(data, state.get("ResultPath", "$"), result)
elif kind == "Map":
items = get_path(data, state.get("ItemsPath", "$"))
result = [run(state["ItemProcessor"], item, tasks, log) for item in items]
data = set_path(data, state.get("ResultPath", "$"), result)
elif kind == "Parallel":
result = [run(b, data, tasks, log) for b in state["Branches"]]
data = set_path(data, state.get("ResultPath", "$"), result)
elif kind == "Succeed":
log.append(f"{name}: SUCCEEDED")
return data
elif kind == "Fail":
log.append(f"{name}: FAILED ({state['Error']})")
raise TaskFailed(state["Error"], state.get("Cause", ""))
except TaskFailed as e:
for c in state.get("Catch", []):
if matches(e.error, c["ErrorEquals"]):
log.append(f"{name}: caught {e.error} -> {c['Next']}")
data = set_path(data, c.get("ResultPath", "$"), {"Error": e.error, "Cause": e.cause})
name = c["Next"]
break
else:
raise
continue
log.append(f"{name}: ok")
if state.get("End"):
return data
name = state["Next"]
import json
machine = json.loads(ASL)
def make_tasks(glue_failures):
state = {"glue_calls": 0}
def glue(params):
state["glue_calls"] += 1
if state["glue_calls"] <= glue_failures:
raise TaskFailed("Glue.ConcurrentRunsExceededException", "another run is active")
return {"JobRunState": "SUCCEEDED", "Arguments": params["Arguments"]}
return {
"arn:aws:states:::lambda:invoke": lambda p: {"Payload": {"files": ["a.json", "b.json"], "valid": True}},
"arn:aws:states:::glue:startJobRun.sync": glue,
"arn:aws:states:::aws-sdk:glue:startCrawler": lambda p: {},
"arn:aws:states:::sns:publish": lambda p: {"MessageId": "m-1"},
}
for failures in (2, 9):
log = []
try:
out = run(machine, {"prefix": "raw/orders/", "run_date": "2026-10-05"}, make_tasks(failures), log)
outcome = "SUCCEEDED, Glue state " + out["results"][0]["JobRunState"]
except TaskFailed as e:
outcome = f"FAILED with {e.error}"
print(f"--- Glue busy for {failures} attempts: {outcome}")
print("\n".join(log))
--- Glue busy for 2 attempts: SUCCEEDED, Glue state SUCCEEDED
ListNewFiles: ok
ValidateFile: ok
ValidateFile: ok
ValidateFiles: ok
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 1
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 2
RunGlueJob: ok
StartCrawler: ok
TransformAndCatalogue: ok
Done: SUCCEEDED
--- Glue busy for 9 attempts: FAILED with PipelineFailed
ListNewFiles: ok
ValidateFile: ok
ValidateFile: ok
ValidateFiles: ok
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 1
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 2
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 3
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 4
RunGlueJob: Glue.ConcurrentRunsExceededException, retry 5
TransformAndCatalogue: caught Glue.ConcurrentRunsExceededException -> NotifyFailure
NotifyFailure: ok
PipelineFailed: FAILED (PipelineFailed)
Retriers are evaluated in order and the first match wins, so put specific errors before States.ALL. Each retrier counts its own attempts. A Catch writes the error and cause into the state data at its ResultPath, so the alert can include them. Set TimeoutSeconds (and HeartbeatSeconds for long callbacks) on tasks, otherwise a stuck task waits until the workflow’s own limit.
Standard workflows can also be redriven: restart a failed execution from the state that failed, keeping the successful steps, instead of rerunning everything.
Pitfalls.
- Retrying non-idempotent steps (an
INSERTwithout a key, sending an email) duplicates side effects. Make steps idempotent or do not retry them. States.ALLretries on everything, including bugs and permission errors that will never succeed. Retry transient errors only.- Catching an error and ending in
Succeed, which hides failures from monitoring.
In interviews. Write a retrier from memory with interval, attempts, backoff and jitter, explain first-match ordering, and show a Catch that records the error and sends an alert. Mention redrive and idempotency.
Integration with Glue and Lambda
What it is. Step Functions calls other services through optimised integrations (with Step Functions-specific behaviour) and AWS SDK integrations (any API, arn:aws:states:::aws-sdk:service:action).
How it works. Three integration patterns:
| Pattern | Resource suffix | Behaviour | Example |
|---|---|---|---|
| Request response | none | Call the API and move on immediately | Start a crawler, publish to SNS |
| Run a job | .sync |
Start the job and wait until it finishes, failing the state if the job fails | glue:startJobRun.sync, athena:startQueryExecution.sync, elasticmapreduce:addStep.sync |
| Wait for callback | .waitForTaskToken |
Pass a task token and pause until something calls SendTaskSuccess or SendTaskFailure |
Human approval, an external system finishing work |
- Lambda:
arn:aws:states:::lambda:invokewraps the response inPayload; select what you need withResultSelector. Add a retrier for Lambda’s transient errors, which AWS recommends for every Lambda task. - Glue:
glue:startJobRun.syncstarts the job, polls until it ends and returns the run details; the state fails if the run fails or times out. Pass arguments such as--run_dateso reruns and backfills are deterministic. - Crawlers have no
.syncintegration. Start them with the SDK integration and, if you must wait, loop withWait,getCrawlerandChoiceuntil the state isREADY.
Pitfalls.
- Using request-response for a Glue job and moving on before it finishes.
- Polling loops without a maximum number of iterations; add a counter or a timeout.
In interviews. Name the three patterns and say which you use for Glue (.sync), for alerts (request response) and for human approval (task token).
Orchestration patterns
Common shapes for data pipelines on Step Functions:
- Sequential with gates: ingest, validate, transform, data quality, publish; each step only if the previous succeeded, with a Catch that alerts.
- Fan-out and fan-in: a Map over files or partitions, or Parallel branches for independent jobs, then a join step.
- Poll until ready:
Wait, check status,Choice, for services without.sync(crawlers, external APIs), with a maximum iteration count. - Callback:
.waitForTaskTokenfor human approval or for an external system that signals completion. - Event-driven start: an EventBridge rule on an S3
Object Createdevent, or EventBridge Scheduler, starts the state machine with the date or object key as input. - Parent and child workflows:
states:startExecution.syncruns a reusable child state machine, keeping each one small and testable. - Compensation (saga): on failure, run undo steps (drop the half-written partition, restore the previous table version) before failing.
- Idempotent backfills: run the same state machine for a list of dates with a Map, each iteration overwriting only its own partition.
| Need | Step Functions | Amazon MWAA (Airflow) |
|---|---|---|
| Serverless, pay per use | Yes | No, an environment runs continuously |
| Python-defined DAGs, rich scheduling and backfill UI | No (JSON or JSONata; CDK can generate) | Yes |
Deep AWS integration with .sync waits |
Yes | Through operators and sensors |
| Multi-cloud and open-source operators | Limited | Yes |
In interviews. Pick two or three patterns and sketch them for a concrete pipeline, then compare with Airflow on the axes in the table.
Practice questions
Why would you use a Standard rather than an Express workflow for a nightly Glue pipeline?
The pipeline runs longer than five minutes (Express’s maximum), needs the .sync integration to wait for Glue (not supported by Express), benefits from exactly-once state execution, and needs a full execution history for debugging and redrive. Express is for short, high-volume, idempotent event processing.
Write a retrier for throttling errors on a Lambda task.
"Retry": [{"ErrorEquals": ["Lambda.TooManyRequestsException", "Lambda.ServiceException"], "IntervalSeconds": 2, "MaxAttempts": 3, "BackoffRate": 2.0, "JitterStrategy": "FULL"}]. It waits about 2, 4 and 8 seconds (randomised by jitter) before giving up; a Catch after it routes to an alert.
You need to process 200,000 files in a prefix. How do you do it with Step Functions?
A Distributed Map with an ItemReader that lists the S3 prefix, an ItemBatcher to group files into batches, Express child executions that call a Lambda or start a job per batch, a MaxConcurrency that protects downstream systems, a tolerated failure percentage, and a ResultWriter to S3. An inline Map would exceed the payload and history limits.
How do you wait for a Glue crawler to finish, since it has no .sync integration?
Start it with the SDK integration (aws-sdk:glue:startCrawler), then loop: Wait 30 seconds, call aws-sdk:glue:getCrawler, and use a Choice on the crawler state; continue when it is READY, with a counter or timeout to stop the loop. Alternatively react to the crawler’s EventBridge event.
What happens to the data when a Catch fires, and why does ResultPath matter?
The Catch passes the state’s input on to its Next state, with the error name and cause inserted at the Catch’s ResultPath (for example $.error). With the default ResultPath of $, the error object replaces the whole input, and later states lose fields such as run_date. Setting ResultPath keeps the input and adds the error.
Key takeaways
- Step Functions runs JSON-defined state machines with Task, Choice, Wait, Parallel, Map, Pass, Succeed and Fail states.
- Standard workflows run up to a year with exactly-once states and full history; Express run up to five minutes at high volume.
- Parallel runs different branches; Map repeats work per item, and Distributed Map scales to 10,000 parallel child executions from S3.
- Retry with interval, backoff, maximum delay and jitter for transient errors; Catch to alert; make steps idempotent.
- Use
.syncfor Glue, Athena and EMR jobs, request response for fire-and-forget calls, and task tokens for callbacks. - State payloads are limited to 256 KiB, so pass references to data in S3, not the data.
Progress is saved in this browser only. No account needed.