skip to content

In Prefect, what does calling a task's .map() method do to the arguments you pass?

level: middleimportance: should knowfreq 55%

answer

  1. One run per element of the input
  2. Two lists pair up, they do not cross
  3. Scalars ride along unchanged
  4. A constant list needs protecting
  5. Each child fails on its own

basics

~20 s

Prefect's .map() creates one task run per element of each iterable argument, zipped element-wise, and runs them through the task runner. Non-iterable values are broadcast to every child run; wrap an iterable you want passed whole in unmapped().

solid answer

~40 s

`process.map(files)` fans out: Prefect creates one task run per element of `files` and dispatches them to the flow's task runner, returning a collection of futures rather than one. With several iterable arguments they are **zipped element-wise**, so `process.map(files, dates)` pairs the first file with the first date — it is not a cartesian product, and mismatched lengths are an error. A plain scalar such as a config string is broadcast unchanged to every child run. The trap is an argument that *is* iterable but should stay whole, like a list of columns or a lookup dict: without `unmapped(columns)` Prefect maps over it too and you get one run per column. You can also map over futures, so `transform.map(extract.map(files))` chains a fan-out. Collect with `.result()` on the returned collection.

code

python · 24 lines
python
from prefect import flow, task, unmapped

COLUMNS = ["id", "ts", "amount"]


@task(retries=2)
def clean(path: str, run_date: str, columns: list[str], bucket: str) -> str:
    return f"{bucket}/{run_date}/{path}"


@task
def summarise(paths: list[str]) -> int:
    return len(paths)


@flow
def ingest(paths: list[str], dates: list[str]) -> int:
    cleaned = clean.map(
        paths,               # mapped: one run per path
        dates,               # mapped: zipped with paths, not crossed
        unmapped(COLUMNS),   # iterable held constant
        "s3://raw",          # scalar broadcast automatically
    )
    return summarise(cleaned)

go deeper

for a junior

Recognise the fan-out shape: mapping a task over a list produces one run per element instead of a loop, and the results come back as a collection you can pass on.

for a middle

Explain argument handling precisely — iterables zip element-wise, scalars broadcast, and a constant iterable must be wrapped so it is not treated as another axis.

for a senior

Show you have operated a wide fan-out: per-element failure and retry, bounding concurrency against a shared database, batching tiny elements, and keeping payloads out of memory.

for a principal

Own the sizing policy — how wide a fan-out a single flow should ever be, when the work belongs in a data-processing engine instead of thousands of orchestrated runs, and what partial success means to consumers.

## What mapping is for Mapping is Prefect's fan-out primitive: the same unit of work applied to every element of a collection, with each element getting its own tracked task run. The typical shapes are one run per file in a landing bucket, per partition, per tenant or per API page. ```python @flow def ingest(files: list[str]): results = clean.map(files) summarise(results) ``` If `files` has 200 entries you get 200 `clean` task runs. Each has its own state, its own logs, its own retries and its own caching, which is exactly why mapping is preferable to a `for` loop inside a single task: when file 137 fails, only that run fails, only that run retries, and the UI shows which one. ## How arguments are interpreted This is the part interviewers probe, because the rule is not obvious. **Iterable arguments are mapped.** Each element becomes the argument for one child run. **Multiple iterables are zipped, not crossed.** ```python clean.map(files, dates) ``` produces one run for `(files[0], dates[0])`, one for `(files[1], dates[1])`, and so on — not `len(files) * len(dates)` runs. If you genuinely want a cross product, build it yourself with `itertools.product` and map over the resulting pairs. Iterables of unequal length are a mistake Prefect will not silently paper over. **Non-iterable arguments are broadcast.** A scalar such as `bucket="s3://raw"` or an `int` is passed unchanged to every child run; you do not need to wrap it. **Iterables you want passed whole need `unmapped()`.** ```python from prefect import unmapped clean.map(files, columns=unmapped(COLUMNS), schema=unmapped(schema_dict)) ``` Without `unmapped`, a list of ten column names would be treated as another axis to map over, and Prefect would either zip it against `files` or raise on the length mismatch. This is the single most common mapping bug: an argument that happens to be a list, a tuple, a dict or a `set` gets mapped when it was meant to be constant. The rule of thumb is that `unmapped` exists for iterable static arguments only. ## Mapping over futures Mapped calls accept futures and the collection returned by a previous `.map()`, so fan-outs chain: ```python extracted = extract.map(files) cleaned = transform.map(extracted) load(cleaned) ``` Prefect wires each `transform` run to the corresponding `extract` run, so element *i* flows through the chain without a barrier that waits for all extracts to finish first. That elementwise wiring is a real advantage over collecting everything and re-mapping. ## Collecting results `.map()` gives you back a collection of futures. You can resolve them individually, or resolve the collection to get a list of values in input order. Passing the collection into a downstream task also works and triggers resolution: `summarise(cleaned)` waits for every child run and receives the list of values — a natural fan-in. ## Failure semantics Each mapped run succeeds or fails independently, and retries configured on the task apply per run, not to the batch. Three consequences worth stating: - **Partial failure is normal.** 198 of 200 runs can succeed. If a downstream task consumes the whole collection, resolving it re-raises the first failure, so the fan-in fails even though most of the work is fine. - **Handle it deliberately** when partial progress is acceptable — resolve futures with failure-tolerant options or inspect states — rather than pretending the batch is atomic. - **Idempotency matters more here.** A retried mapped run re-executes its element, so writing must be safe to repeat for that element. ## Sizing the fan-out Mapping thousands of elements is where fan-out stops being free. Every child run is a tracked object with state transitions and logs, and the flow process holds the futures. Practical guidance: - Batch the input so one run handles a chunk of files rather than one file, when per-item work is tiny. - Bound the parallelism through the task runner's worker count, or with tag-based concurrency limits, so 500 mapped runs do not open 500 warehouse connections. - Watch the flow process's memory: results held in futures live somewhere, and mapping over huge in-memory payloads is a footgun — pass keys or paths, not dataframes. ## The interview answer Say it in this order: one run per element; several iterables zip; scalars broadcast; iterables you want constant need `unmapped`; the result is a collection of futures you can chain or fan back in. Then add the operational nuance about per-run failure and bounding the width, which is what separates someone who has read the docs from someone who has run a 2,000-way fan-out at 2am.

  • Why does a mapped task run per element beat a for loop inside one task?
    Granularity. Each element gets its own task run, so it retries, caches, times out and reports independently, and the failing element is visible in the UI. A loop inside one task retries the whole batch, redoing the 137 files that already succeeded, and shows a single opaque failure.
  • How would you produce a cartesian product across two lists with a mapped task?
    Build the pairs yourself before mapping, for example with `itertools.product`, then map over the resulting sequence of tuples or over two aligned lists derived from it. Mapping several iterables directly zips them element-wise, so it will never generate a cross product for you.
  • What stops a 500-way fan-out from overwhelming a downstream database?
    Bound the concurrency rather than the fan-out: cap the task runner's worker count, or apply a tag-based concurrency limit so only N of those runs execute at once. Also consider batching elements so each run handles a chunk, which cuts the number of connections and the orchestration overhead.

saying these in an interview costs you the question

  • Expecting a cartesian product from two mapped iterables
  • Forgetting unmapped and mapping over a constant list
  • Treating a mapped batch as atomic all-or-nothing
  • Mapping ten thousand tiny elements instead of batching
  • Passing large in-memory objects into every child run

context