Repository navigation
DateTimeSensorAsync support start_from_trigger and templated target_time - #73928
raphaelauv wants to merge 2 commits into
Conversation
79d3512 to
06ac36a
Compare
06ac36a to
d2115cd
Compare
dabla
left a comment
There was a problem hiding this comment.
Thanks for picking this up. The __init__ part is right: passing the raw target_time in trigger_kwargs is how the triggerer-side rendering from #55068 is meant to be used, and the Airflow < 3.3 fallback is kept.
Two things need fixing before merge:
execute()is implemented the wrong way. It is the worker path and should checkself.start_from_trigger, not the Airflow version, and keep building the trigger frommoment=self._moment. That limitstarget_timeto the triggerer path and makes the task-sdk test changes unnecessary.- On Airflow 3.3+ a
datetimetarget_time(with or withoutstart_from_trigger) now reachesDateTimeTriggerastarget_timeand fails in the triggerer with aTypeError, where it works onmain. I reproduced this against the branch. HandlingdatetimeinDateTimeTrigger.momentfixes it.
The tests need some work as well: DateTimeTrigger has no tests for the new target_time argument, and the new end-to-end test combines the worker path and the triggerer rendering in a way that does not happen at run time.
The failing "Compat 3.0.6" job looks unrelated: it timed out in other providers' tests before reaching the standard provider.
|
|
||
| if AIRFLOW_V_3_3_PLUS: | ||
| trigger = DateTimeTrigger( | ||
| target_time=self.target_time, |
There was a problem hiding this comment.
[blocker] execute() branches on the wrong condition: it should check self.start_from_trigger, not the Airflow version, and keep building the trigger from moment.
There are two separate paths:
- Triggerer path (
start_from_trigger=True): the scheduler defers the task directly fromstart_trigger_args, and the triggerer renders the rawtarget_timefromtrigger_kwargs.execute()is not involved. - Worker path (
start_from_trigger=False): the worker has already renderedself.target_timewhenexecute()runs, so the trigger can be built frommoment=self._momentexactly as onmain.
Selecting target_time= in execute() on 3.3+ applies the triggerer-path mechanism to the worker path, which has two costs:
- It ships the rendered value with whatever type it has. A
datetime(passed directly, or rendered natively) reaches the trigger astarget_timeand fails there, see my comment onDateTimeTrigger.moment. - Plain deferral now depends on the triggerer running this provider version. An older
DateTimeTriggerrejects the unknowntarget_timekwarg, whilemomentworks with every version.
Suggestion:
def execute(self, context: Context) -> None:
if not self.start_from_trigger:
if AIRFLOW_V_3_0_PLUS:
trigger = DateTimeTrigger(moment=self._moment, end_from_trigger=self.end_from_trigger)
else:
trigger = DateTimeTrigger(moment=self._moment)
self.defer(method_name="execute_complete", trigger=trigger)With that, the changes to test_supervisor.py and test_task_runner.py are no longer needed.
| """ | ||
|
|
||
| def __init__(self, moment: datetime.datetime, *, end_from_trigger: bool = False) -> None: | ||
| template_fields = ("target_time",) |
There was a problem hiding this comment.
[nit] This class attribute has no effect, and the docstring above describes a mechanism that is not the one in use.
BaseTrigger.__init__ sets self.template_fields = () on the instance, which shadows the class attribute:
>>> DateTimeTrigger(target_time="{{ ds }}").template_fields
()The fields that get rendered are set by the task_instance setter in BaseTrigger: the operator's template_fields that are also keys of start_trigger_args.trigger_kwargs and attributes of the trigger. So rendering works here because DateTimeSensor.template_fields contains target_time and the sensor puts target_time in trigger_kwargs. I would drop this line and reword the docstring accordingly.
| def moment(self) -> pendulum.DateTime: | ||
| if self._moment is None: | ||
| # Resolved lazily: by now the triggerer has rendered target_time in place. | ||
| if not isinstance(self.target_time, str) or not self.target_time: |
There was a problem hiding this comment.
[blocker] A datetime target_time is accepted by __init__ but makes the trigger fail at run time.
The signature allows target_time: datetime.datetime | str | None, but moment only resolves strings. Reproduced against this branch:
trigger = DateTimeTrigger(target_time=datetime.datetime(2020, 1, 1, tzinfo=datetime.timezone.utc))
await trigger.run().__anext__()
# TypeError: DateTimeTrigger has neither a 'moment' nor a usable 'target_time'On Airflow 3.3+ the sensor now sends a datetime this way in cases that work on main:
DateTimeSensorAsync(target_time=<datetime>, start_from_trigger=True):__init__puts the raw value intrigger_kwargs. The updatedtest_async_start_from_trigger_localizes_naive_datetimeasserts exactly that (a naivedatetimeundertarget_time).DateTimeSensorAsync(target_time=<datetime>)throughexecute(): the updated expectation intest_supervisor.pyshows adatetime.datetimebeing sent astarget_time.- A template that renders to a native
datetime(render_template_as_native_obj=True) would end up in the same place.
DateTimeSensor._moment already handles both types, so the trigger can mirror it:
@property
def moment(self) -> pendulum.DateTime:
if self._moment is None:
target_time = self.target_time
if isinstance(target_time, datetime.datetime):
target_time = target_time.isoformat()
if not isinstance(target_time, str) or not target_time:
raise TypeError(f"Expected str or datetime.datetime type for target_time. Got {type(target_time)}")
self._moment = timezone.convert_to_utc(timezone.parse(target_time))
return self._momenttimezone.parse also localizes naive values, which restores the behaviour the test name promises. Please add trigger tests in triggers/test_temporal.py for target_time as str, aware datetime and naive datetime, for the "exactly one of" validation, and for the serialize() round trip.
|
|
||
| @pytest.mark.asyncio | ||
| @pytest.mark.skipif(not AIRFLOW_V_3_3_PLUS, reason="Test only for AF < 3.2") | ||
| async def test_full_run_worker_path_templated_past_target_time(self): |
There was a problem hiding this comment.
[warning] This test exercises a flow that does not occur in production, so it does not cover the new feature.
It calls op.execute(ctx) without rendering the operator first (the worker always renders before execute), then renders the template on the trigger through a SimpleNamespace task instance. With start_from_trigger=True the scheduler defers the task directly and execute() is never called; without it, the triggerer does not render anything.
A test closer to the real path would build the trigger the way the triggerer does:
trigger = DateTimeTrigger(**op.start_trigger_args.trigger_kwargs)
trigger.task_instance = ... # object exposing task_id and task=op
trigger.render_template_fields(ctx)
event = await trigger.run().__anext__()Smaller points in this file:
- The
skipifreasons do not match the conditions:not AIRFLOW_V_3_3_PLUSwith "Test only for AF < 3.2" skips on < 3.3, and theAIRFLOW_V_3_3_PLUSone below should read "< 3.3". TaskDeferredandSimpleNamespaceshould be imported at the top of the module.test_async_start_from_trigger_molocalmentlooks like an accidental rename.- The new
UserWarningon Airflow < 3.3 has no test (pytest.warns).
| assert op.start_trigger_args.trigger_kwargs["moment"] == pendulum.datetime(2020, 1, 1, tz="UTC") | ||
|
|
||
| if AIRFLOW_V_3_3_PLUS: | ||
| assert op.start_trigger_args.trigger_kwargs["target_time"] == datetime.datetime(2020, 1, 1, 0, 0) |
There was a problem hiding this comment.
[warning] This assertion pins the bug described on DateTimeTrigger.moment.
The test is named ..._localizes_naive_datetime, but on 3.3+ it now asserts that the naive datetime is passed through unchanged. The trigger then rejects it at run time. Once the trigger handles datetime, I would extend this test to check the resolved value, for example DateTimeTrigger(**op.start_trigger_args.trigger_kwargs).moment == pendulum.datetime(2020, 1, 1, tz="UTC").
|
|
||
| trigger_kwargs_key_name = "moment" | ||
| trigger_kwargs_class_name = "pendulum.datetime.DateTime" | ||
| if AIRFLOW_V_3_3_PLUS: |
There was a problem hiding this comment.
[warning] A change in the standard provider should not need edits in the task-sdk tests.
These two tests only use DateTimeSensorAsync as a sample deferrable task and assert the trigger kwargs it sends. Their expectations had to change because execute() now sends target_time instead of moment. Once execute() builds the trigger from moment again (see my comment there), both test_supervisor.py and test_task_runner.py can be reverted.
The AIRFLOW_V_3_3_PLUS branches added here are also dead code: the task-sdk tests only run against the current main (3.4.0), never against older Airflow versions.
f3fbaac to
84271f8
Compare
following #72659
since airflow 3.3 start_from_trigger works with templated values
Was generative AI tooling used to co-author this PR?