Skip to content

Commit 8a44bf7

Browse files
committed
Say when not to iterate an operator that finds its remote work by the task instance's identity
1 parent 616fea8 commit 8a44bf7

2 files changed

Lines changed: 26 additions & 0 deletions

File tree

‎task-sdk/docs/mapped-tasks-vs-iterable-tasks.rst‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -463,6 +463,15 @@ Avoid Iterable Tasks when:
463463
- Sub-tasks need to defer (deferrable operators) or reschedule (reschedule-mode sensors) — a
464464
sub-task index has no task instance of its own to defer or reschedule against, so either raises
465465
a non-retryable failure instead of pausing.
466+
- The operator finds its remote work by the task instance's identity. Every item runs as the one
467+
iterated task instance, with its ``dag_id``, ``task_id``, ``run_id`` and ``map_index``: a
468+
``KubernetesPodOperator`` that reattaches (``durable``, or ``reattach_on_restart`` before it)
469+
looks a running pod up by those labels and adopts a sibling item's pod, or fails on finding
470+
several; an ``EcsRunTaskOperator`` with ``reattach=True`` builds its ``startedBy`` from the same
471+
fields. Where the operator takes labels of its own, one that names the item tells the jobs apart:
472+
``labels={"index": "{{ ti.index }}"}`` on the ``KubernetesPodOperator``, as every template of an
473+
item renders against the item's own ``ti``. Otherwise turn reattachment off, or map the task with
474+
``.expand()``.
466475

467476
.. tip::
468477

‎task-sdk/tests/task_sdk/definitions/test_iterableoperator.py‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1195,6 +1195,23 @@ def test_a_template_reads_the_indexed_tasks_own_state_store(self):
11951195

11961196
assert task.task.offset == "7"
11971197

1198+
def test_a_template_reads_the_items_own_index(self):
1199+
"""
1200+
``{{ ti.index }}`` in a partial kwarg renders per item, so an operator that finds its remote
1201+
work by labels (a ``KubernetesPodOperator`` that reattaches) can be told the item apart.
1202+
"""
1203+
with DAG("test_dag") as dag:
1204+
expand_input = ListOfDictsExpandInput([{"arg1": "a"}, {"arg1": "b"}])
1205+
mapped_op = MockOperator.partial(task_id="indexed_label", dag=dag, arg2="{{ ti.index }}")._expand(
1206+
expand_input, strict=True, register_with_dag=False
1207+
)
1208+
iterable_op = IterableOperator(operator=mapped_op, expand_input=expand_input, dag=dag)
1209+
1210+
with mock_context(task=iterable_op) as context:
1211+
results = sorted(iterable_op.execute(context=context))
1212+
1213+
assert results == [("a", "0", None), ("b", "1", None)]
1214+
11981215
def test_execute_renders_template_fields_off_the_loop_thread(self):
11991216
"""
12001217
Rendering a sub-task's template fields may call the supervisor synchronously: an XComArg in

0 commit comments

Comments
 (0)