Menu

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.

  • Intermediate
  • 19 min read
  • Updated Oct 2026
On this page
  1. Sample pipeline
  2. State machines
  3. Standard vs Express
  4. Map and Parallel states
  5. Error handling and retries
  6. Integration with Glue and Lambda
  7. Orchestration patterns
  8. Practice questions
  9. Key takeaways

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 lose run_date.
  • An execution role without permission for each service the definition calls (including glue:GetJobRun for .sync polling and events: 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. MaxConcurrency limits 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 an ItemReader (a list of objects, a CSV or JSON file, or an S3 Inventory manifest), group items with ItemBatcher, write results to S3 with ResultWriter, 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 ToleratedFailurePercentage or 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 INSERT without a key, sending an email) duplicates side effects. Make steps idempotent or do not retry them.
  • States.ALL retries 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:invoke wraps the response in Payload; select what you need with ResultSelector. Add a retrier for Lambda’s transient errors, which AWS recommends for every Lambda task.
  • Glue: glue:startJobRun.sync starts 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_date so reruns and backfills are deterministic.
  • Crawlers have no .sync integration. Start them with the SDK integration and, if you must wait, loop with Wait, getCrawler and Choice until the state is READY.

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:

  1. Sequential with gates: ingest, validate, transform, data quality, publish; each step only if the previous succeeded, with a Catch that alerts.
  2. Fan-out and fan-in: a Map over files or partitions, or Parallel branches for independent jobs, then a join step.
  3. Poll until ready: Wait, check status, Choice, for services without .sync (crawlers, external APIs), with a maximum iteration count.
  4. Callback: .waitForTaskToken for human approval or for an external system that signals completion.
  5. Event-driven start: an EventBridge rule on an S3 Object Created event, or EventBridge Scheduler, starts the state machine with the date or object key as input.
  6. Parent and child workflows: states:startExecution.sync runs a reusable child state machine, keeping each one small and testable.
  7. Compensation (saga): on failure, run undo steps (drop the half-written partition, restore the previous table version) before failing.
  8. 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 .sync for 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.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Workflow types and quotas checked against the AWS Step Functions Developer Guide in October 2026. The state machine definition was checked to parse as JSON and was run through a small local Python interpreter (Python 3.11) with mocked tasks; it models Task, Map, Parallel, Retry and Catch only and is not the Step Functions service. CLI commands were written from the documentation and not executed.

Progress is saved in this browser only. No account needed.

Search
Filter by type