skip to content

A Step Functions workflow must process several million JSON objects stored in an S3 bucket. Why would you use a Distributed Map state instead of an inline Map?

level: seniorimportance: should knowfreq 44%

answer

  1. one loops inside the parent execution
  2. payload size gates the item list
  3. the parent history has an event ceiling
  4. children carry their own histories
  5. results land in S3, not the payload

basics

~20 s

An inline Map runs every iteration inside the parent execution, so the items must fit in the state payload and every iteration's events land in one execution history, which is bounded. Distributed Map reads items straight from S3 and runs them as separate child executions at far higher concurrency.

solid answer

~50 s

An inline `Map` iterates over an array that is already in the state's payload, and each iteration's events are recorded in the *parent* execution's history. Millions of items break both: the payload cannot hold the array, and a Standard execution's history has a hard event ceiling you would blow through long before finishing. Concurrency is also modest. **Distributed Map** changes the execution model: an `ItemReader` reads the dataset directly from S3 — a JSON array, JSON Lines, CSV, a bucket listing, or an S3 Inventory manifest — `ItemBatcher` groups items per child, and each batch runs as its own **child execution** (usually Express) with its own history, up to a very high parallelism. `ToleratedFailurePercentage` lets the parent survive a few bad records instead of aborting everything, and `ResultWriter` writes results back to S3 rather than trying to return them in the payload. That is the difference between a fan-out helper and a batch-processing engine.

code

json · 33 lines
json
{
  "Comment": "Distributed Map over an S3 prefix",
  "StartAt": "ProcessObjects",
  "States": {
    "ProcessObjects": {
      "Type": "Map",
      "ItemReader": {
        "Resource": "arn:aws:states:::s3:listObjectsV2",
        "Parameters": { "Bucket": "raw-events", "Prefix": "2025/08/" }
      },
      "ItemBatcher": { "MaxItemsPerBatch": 100 },
      "MaxConcurrency": 500,
      "ToleratedFailurePercentage": 1,
      "ItemProcessor": {
        "ProcessorConfig": { "Mode": "DISTRIBUTED", "ExecutionType": "EXPRESS" },
        "StartAt": "HandleBatch",
        "States": {
          "HandleBatch": {
            "Type": "Task",
            "Resource": "arn:aws:states:::lambda:invoke",
            "Parameters": { "FunctionName": "handle-batch", "Payload.$": "$" },
            "End": true
          }
        }
      },
      "ResultWriter": {
        "Resource": "arn:aws:states:::s3:putObject",
        "Parameters": { "Bucket": "run-results", "Prefix": "august" }
      },
      "End": true
    }
  }
}

go deeper

for a junior

Know that Map repeats one branch per item in an array, and that the array normally comes from the state's input via ItemsPath. Recognising that huge datasets need a different mode is enough here.

for a middle

Explain the two concrete limits that break inline Map — the state payload cap and the parent execution's event ceiling — and describe how child executions relocate that work.

for a senior

Show the operational reasoning: batch size against cost and retry granularity, MaxConcurrency chosen from the weakest downstream dependency, tolerated-failure thresholds, and results written to S3.

for a principal

Own the build-versus-use call — when Distributed Map is the right batch engine versus a Glue or EMR job, and what the fan-out does to concurrency budgets and spend across the whole account.

## Two states with the same name and different physics `Map` has two processing modes, selected by `ItemProcessor.ProcessorConfig.Mode`: `INLINE` (the original behaviour, and the default) and `DISTRIBUTED`. They look almost identical in the definition and behave nothing alike at scale. ## Why inline Map stops working An inline Map is a loop *inside* one execution: 1. **The array must already be in the payload.** `ItemsPath` selects an array out of the state input, and state payloads are capped at 256 KiB (as of 2025). A million object keys do not fit, so you cannot even express the input. 2. **Every iteration writes into the parent's execution history.** A Standard execution's history is bounded (25,000 events as of 2025), and each iteration contributes several events. Even a few thousand items can exhaust it, and when the history limit is hit the execution fails outright — after having done most of the work. 3. **Concurrency is modest.** Inline Map's `MaxConcurrency` is capped in the low tens (40 as of 2025). At millions of items that is not throughput, it is a queue. 4. **All-or-nothing failure.** One iteration failing without handling takes down the Map, and with it the run. Inline Map is right for tens or a few hundred items already in hand — split an order into line items, fan out to a handful of regions. ## What Distributed Map changes Distributed Map turns the state into a coordinator of **child executions**: ```json "ProcessObjects": { "Type": "Map", "ItemReader": { "Resource": "arn:aws:states:::s3:listObjectsV2", "Parameters": { "Bucket": "raw-events", "Prefix": "2025/08/" } }, "ItemBatcher": { "MaxItemsPerBatch": 100 }, "MaxConcurrency": 1000, "ToleratedFailurePercentage": 1, "ItemProcessor": { "ProcessorConfig": { "Mode": "DISTRIBUTED", "ExecutionType": "EXPRESS" }, "StartAt": "Handle", "States": { "Handle": { "Type": "Task", "Resource": "arn:aws:states:::lambda:invoke", "Parameters": { "FunctionName": "handle-batch" }, "End": true } } }, "ResultWriter": { "Resource": "arn:aws:states:::s3:putObject", "Parameters": { "Bucket": "run-results", "Prefix": "august" } }, "End": true } ``` Four mechanisms carry the weight: - **`ItemReader`** reads the dataset where it lives instead of through the payload. It supports a JSON array in an object, JSON Lines, CSV (with a header row or supplied headers), an S3 bucket listing, and an S3 Inventory manifest. The payload limit stops being relevant because the items never travel through it. - **`ItemBatcher`** groups items so each child processes many records — the difference between one child per object (expensive, chatty) and one child per hundred objects. `BatchInput` lets you attach constant context to every batch. - **Child executions** each get their own execution history and their own limits, so the parent's history stays small — it records the map run, not every item. Distributed Map scales to a very high number of concurrent child executions (up to 10,000 as of 2025), far past inline Map's ceiling. `ExecutionType` picks Standard or Express for the children; Express is the usual choice for short per-item work, and it changes the cost shape from per-transition to per-duration. - **`ToleratedFailurePercentage` / `ToleratedFailureCount`** make partial failure a first-class outcome. In a million-record batch, a handful of malformed rows should not abort the run — you set a threshold, and the failed items are recorded so you can reprocess just those. `ResultWriter` closes the loop: results go to S3, because returning a million results in the state output would run straight back into the payload limit you just escaped. ## The costs you take on Distributed Map is not free judgment-wise. Child executions are billed as executions, so batching size directly drives cost. Downstream systems now see up to thousands of concurrent workers — the Lambda function's reserved concurrency, the database's connection budget, or the third-party API's rate limit becomes the real bottleneck, and `MaxConcurrency` is the throttle you set deliberately rather than discovering under load. The parent's execution page also no longer shows per-item detail; you inspect a child execution or the results written to S3. ## How to answer Name the two blockers (payload size and parent execution history), then the model change (child executions with their own histories, reading from S3), then the operational levers (batching, `MaxConcurrency` chosen for the downstream system, tolerated failure, results to S3). That progression — limit, mechanism, operations — is what marks a senior answer over a feature recital.

  • Why does batching matter so much in a Distributed Map?
    Each batch becomes a child execution, and executions are billed and rate-limited as units. One item per child at a million items means a million executions and a million cold-start-shaped invocations; a hundred items per child cuts that by two orders of magnitude. You trade retry granularity for cost — a failed batch reprocesses all its items.
  • What sets the right MaxConcurrency for a Distributed Map?
    The weakest downstream dependency, not the state machine. If each child calls a database, the connection budget caps you; if it calls a partner API, their rate limit does; if it invokes a Lambda, that function's concurrency does. Pick the number from that constraint and set it explicitly, because the default parallelism is high enough to take the dependency down.
  • How do you handle a small number of malformed records in a multi-million-item run?
    Set `ToleratedFailurePercentage` or `ToleratedFailureCount` so the parent completes instead of aborting on the first bad row, and use `ResultWriter` so the per-item outcomes land in S3. You then reprocess just the failed items from that output rather than replaying the whole dataset.

saying these in an interview costs you the question

  • Assuming inline Map just scales to any array size
  • Passing millions of S3 keys through the state payload
  • Leaving MaxConcurrency unset against a fragile dependency
  • Expecting the parent execution history to show every item
  • Returning all per-item results in the state output

context