Uh oh!
There was an error while loading. Please reload this page.
Prevent using trigger_rule=TriggerRule.ALWAYS in a task-generated mapping within mapped task groups - #43368
Conversation
1ef1aee to
ada1eddComparetrigger_rule="always" in a dynamic mapped tasktrigger_rule=TriggerRule.ALWAYS in a dynamic mapped taskpotiuk
commented
Oct 24, 2024
cc: @uranusjr -> I think it will need your input. |
I need some help with fixing the tests - it gives strange errors when running them altogether: Thanks beforehand :) Update: I've managed by moving the |
1bd0fe9 to
22f6bd0Compare22f6bd0 to
c753ca2Comparepotiuk
commented
Oct 31, 2024
Can you resolve conflicts @shahar1 and back-port it - we are going to have 2.10.3 RC2 so we might include it there |
c753ca2 to
e057029Comparejedcunningham
commented
Dec 4, 2024
So it does work if you use static list to expand: @task(trigger_rule="always")defhello(input):
print(f"Hello, {input}")
hello.expand(input=["world", "moon"])Granted, not sure it matters. I question the value of "always" anyway - there is no point in having any upstream if it's just going to run right away anyways, right? Just use a root task... |
It does matter if upstream is |
That's correct - it only applies for dynamic task mappings, where the |
Sorry, I should have spent more time digging into this - I leaned too much on the title/newsfrag/docs changes. I will say, however, that the docs/newsfragment on this is misleading/wrong. It does still work for static expansions. This also only blocks task group expansions, expanding bare tasks with dynamic input from another task still parses and ultimately fails when running: @taskdefget_input():
return ["world", "moon"]
@task(trigger_rule="always")defhello(input):
print(f"Hello, {input}")
hello.expand(input=get_input())I'd be in favor of removing this from 2.10.4 for now, then go in and block both bare tasks and task groups, and document it more accurately too. (To be clear, just saying we pull it based on timing, if we can get it before we merge the 2.10.4 changes I'm completely happy with that too.) |
jedcunningham
commented
Dec 6, 2024
@eladkal I'm curious, can you provide a simple example for the short circuit scenario you are talking about? In my experimenting, that feels like a bug on short circuit, not a useful feature of always. e.g. if I have this: @task.short_circuit()defshort():
sleep(20)
returnFalse@task(trigger_rule="always")defhello(input):
print(f"Hello, {input}")
short() >>hello.expand(input=["world", "moon"])(Note: I added a sleep to expand the race condition a bit to make it easier to see the problem.) When it runs, the ALWAYS means those tasks run right away. So those hello tasks have run before the short task finishes: Then, when short eventually finishes, it goes in and retroactively skips... but thats wrong, those have already run! |
eladkal
commented
Dec 6, 2024
then this is a bug. I assume you hit this because you set direct dependency between them. I think if there will be a middle task you won't hit this bug. I do agree that the value of always is very minor. We can consider removing it for Airflow 3. |
trigger_rule=TriggerRule.ALWAYS in a dynamic mapped tasktrigger_rule=TriggerRule.ALWAYS in a task-generated mapping within mapped tasktrigger_rule=TriggerRule.ALWAYS in a task-generated mapping within mapped tasktrigger_rule=TriggerRule.ALWAYS in a task-generated mapping within mapped task groups


closes: #30334
I started to stabilize the dynamic task mapping area, and specifically the interactions with the different trigger rules.
In the case of
TriggerRule.ALWAYS, it is currently impractical to use it along with the dynamic task mapping. That is due the fact that the expanded task is triggered immediately, while the expanded parameters from the upstream are undefined - what will always result in the immediate failure of the task (upstream_failure).I'm not entirely sure what the purpose of the original issue was by putting an "ALWAYS" task within a dynamic task group (they mentioned using a "watcher" pattern). However, I think that for now it makes sense to assume that upon an expension of a task - its expanded parameter(s) should be well-defined.
If at some point we find a good reason to make this definition more flexible - I'm up to it, but at the current implementation it's rather a misuse that we should prevent.
I will cherry-pick into
v2-10-testupon approval.^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named
{pr_number}.significant.rstor{issue_number}.significant.rst, in newsfragments.