Airflow dynamic task mapping: one Sume job per prompt, capped at 4

Use submit.expand(prompt=PROMPTS) with max_active_tis_per_dag=4 and a key built from run_id and map_index. Mind max_map_length (default 1024) for big batches.

4 min readSume
All posts

Airflow's dynamic task mapping turns a list of prompts into one task instance per prompt at run time, and max_active_tis_per_dag caps how many of them run at once. Set the cap to your Sume concurrency, send an Idempotency-Key built from the run id and the map index, and an Airflow retry of a mapped copy becomes a replay instead of a second paid job.

The cap matters because Sume admits a limited number of generations at once per plan and queues a limited number more; submitting a thousand at the same moment is how you reach queue_full.

What the Airflow page constrains

Airflow's documentation says expand() accepts only keyword arguments, that the number of copies is limited by the [core] max_map_length setting with a default of 1024, and that max_active_tis_per_dag limits parallel mapped copies across all active runs of the DAG, not only the current one. The last point is easy to miss: two overlapping runs share the same cap of four, which is what you want when the cap is a vendor limit.

On the Sume side, the generation admission guide gives an example plan with concurrency 4, 20 queued and 24 accepted. A cap of 4 matches it. A larger batch than the queue holds returns 429 queue_full, which the task's retries then absorb, with the same key.

Where each limit lives (read 2026-10-03)
LimitSet byValue in this DAG
Mapped copies per taskcore max_map_length1024 default; chunk larger batches
Parallel copiesmax_active_tis_per_dag4, across all active runs
Retries per copyTask retries argument3, with retry_delay
Duplicate protectionIdempotency-Keyrun_id plus map_index

The DAG

Each mapped copy submits in async mode and returns a job id, which downstream tasks can wait on or collect. The key uses get_current_context() so a retry of the same copy sends the same key.

import os, requests, pendulum
from airflow.decorators import dag, task
from airflow.operators.python import get_current_context

PROMPTS = ["A paper boat on a river", "A lighthouse in fog", "A desert train at dusk"]

@dag(schedule=None, start_date=pendulum.datetime(2026, 10, 1), catchup=False)
def sume_batch():
    @task(max_active_tis_per_dag=4, retries=3)
    def submit(prompt: str) -> str:
        ctx = get_current_context()
        key = f"{ctx['run_id']}-{ctx['ti'].map_index}"
        r = requests.post(
            "https://api.sume.com/v1/images",
            headers={"Authorization": f"Bearer {os.environ['SUME_API_KEY']}",
                     "Idempotency-Key": key},
            json={"model": "sume/auto", "prompt": prompt, "mode": "async"},
            timeout=20,
        )
        r.raise_for_status()  # 429/5xx -> Airflow retry, same key
        return r.json()["data"]["job"]["id"]

    submit.expand(prompt=PROMPTS)

sume_batch()

Batches larger than the map limit

If the list can exceed the map length, split it before mapping, or map over chunks and loop inside each copy with one key per item. Per-item mapping gives the clearest retry boundaries; per-chunk mapping uses fewer task instances but then a retry replays the whole chunk, which is safe only because every item inside it has its own key. Either way, fetch results with a webhook or the jobs status route instead of holding a worker slot open for generation.

Sources

Related posts

More in Developers

All Developers posts

Written by Sume