Repository navigation
Pass the scheduler's session to set_state when processing executor events - #73213
namanjain24-sudo wants to merge 1 commit into
Conversation
c79adcb to
58d3e70
Compare
…ents process_executor_events called ti.set_state() without a session when a cleared task instance is reported terminated and when the Dag of a finished task instance cannot be loaded. set_state is @provide_session and settings.Session is scoped, so create_session() returned the scheduler's own session and committed and closed it on exit, in the middle of the batch.
58d3e70 to
b9d0f3f
Compare
|
Hello @namanjain24-sudo - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 7 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
SchedulerJobRunner.process_executor_eventscallsti.set_state()withoutsessionin two places: when a cleared (RESTARTING) task instance is reported as successfully terminated, and when the Dag of a finished task instance cannot be loaded.TaskInstance.set_stateis@provide_sessionandsettings.Sessionis scoped, socreate_session()returns the scheduler's own session and commits and closes it on exit. This is the same mechanism as #67850 and #71968.process_executor_eventsruns insidewith create_session()in_run_scheduler_loop, not underprohibit_commit, so nothing raises. Instead, in the middle of an executor-event batch:FOR UPDATE SKIP LOCKEDrow locks taken on the batch's task instances;close()detaches the task instances loaded for the batch, so changes made to them afterwards are never written. In a local run with oneRESTARTING→SUCCESSevent and fourQUEUEDevents carrying an external executor id, none of the fourexternal_executor_idvalues reached the database without this change, and all four did with it (SQLite and PostgreSQL 16).This passes the scheduler's session at both call sites, as
_enqueue_task_instances_with_queued_statealready does for its ownti.set_state()call.Tests:
test_process_executor_events_sets_state_in_callers_transactioncovers both paths: the state change has to roll back with the caller's transaction. On main both cases fail (the task instance is already committed asNone/failed); with this change they pass.test_scheduler_job.pyand prek, including mypy for airflow-core, pass locally on SQLite and PostgreSQL 16.Not covered here. In the same loop,
executor.send_callback()reachesDatabaseCallbackSink.send, which is also@provide_sessionand gets no session. So when the task instance has anon_failure_callback/on_retry_callback(reproduced locally with this change applied) or email configured, the scheduler's transaction is still committed at that call._maybe_requeue_stuck_tiand_purge_task_instances_without_heartbeatscallsend_callback()the same way. Fixing that means either adding asessionargument toBaseExecutor.send_callbackandBaseCallbackSink.send, which executors that overridesend_callbackwould then have to accept, or handing the session to the sink directly, asTriggeralready does withDatabaseCallbackSink().send(callback=request, session=session). I kept this PR to theset_statecalls and am happy to follow up with whichever approach you prefer.Was generative AI tooling used to co-author this PR?
Generated-by: a Gen-AI coding assistant, following the guidelines. I reviewed the change and ran the checks above.