Dynamic Task Mapping

Dynamic Task Mapping with expand()

Dynamic task mapping creates task instances at run time, one per element of a list or upstream XCom. partial() fixes the constant arguments, several expand() arguments form a cross product, expand_kwargs() takes one dictionary per instance, and a downstream task receives every result as a lazy sequence:

dags/mapping_demo.py (excerpt): expand, partial and expand_kwargsPython
    @task
    def line_total(qty: int, unit_price: float, discount: float) -> float:
        return round(qty * unit_price * (1 - discount), 2)
...
    crossed = line_total.partial(discount=0.1).expand(qty=[1, 2], unit_price=[14.99, 39.5])
    paired = line_total.expand_kwargs([{"qty": 3, "unit_price": 18.0, "discount": 0.0},
                                       {"qty": 1, "unit_price": 22.5, "discount": 0.2}])
    show("expand (2 x 2):", crossed)
    show("expand_kwargs:", paired)
Output
$ python dags/mapping_demo.py          # the file ends with dag.test() (Section 7.7.5)
expand (2 x 2): [13.49, 35.55, 26.98, 71.1]
expand_kwargs: [54.0, 18.0]

Each instance has a map_index and its own state, retries and log. Classic operators map too (BashOperator.partial(task_id="wc").expand(bash_command=[...])), as do task groups. An empty list skips the task; one longer than [core] max_map_length (1,024) fails it. Map over files, not rows: each instance costs a task's scheduling overhead.