FrankYang0529 opened a new pull request, #73832:
URL: https://github.com/apache/airflow/pull/73832
## Why
- #72048 added `@provide_session` to
`EdgeExecutor.try_adopt_task_instances()` so that the method could read the
`edge_job` table.
- The scheduler calls this method from `adopt_or_reset_orphaned_tasks()`
without passing a session. At that point, the scheduler's own scoped session
still holds the orphaned task instances, which the scheduler loaded with
`load_only(...)`. `@provide_session` gets its session from `create_session()`,
which returns the scheduler's own scoped session. When
`try_adopt_task_instances()` returns, `create_session()` commits and closes the
scheduler's session.
- The scheduler then reads `ti.last_heartbeat_at` on a detached task
instance and exits with `DetachedInstanceError`.
## How
- Read the `edge_job` rows through `create_session(scoped=False)`, which
opens a separate session and leaves the scheduler's session and the orphaned
task instances alone.
## Verification
- Unit test: `uv run --frozen --project providers/edge3 pytest
providers/edge3/tests/unit/edge3/executors/test_edge_executor.py`
- Integration test:
```bash
export AIRFLOW_HOME="${TMPDIR:-/tmp}/edge-adopt-demo"
export AIRFLOW__CORE__DAGS_FOLDER="$AIRFLOW_HOME/dags"
export AIRFLOW__CORE__LOAD_EXAMPLES=False
export AIRFLOW__CORE__PARALLELISM=3
export AIRFLOW__CORE__EXECUTOR=airflow.providers.edge3.executors.EdgeExecutor
export
AIRFLOW__DATABASE__EXTERNAL_DB_MANAGERS=airflow.providers.edge3.models.db.EdgeDBManager
rm -rf "$AIRFLOW_HOME" && mkdir -p "$AIRFLOW__CORE__DAGS_FOLDER"
cat > "$AIRFLOW__CORE__DAGS_FOLDER/edge_adopt_demo.py" <<'EOF'
import pendulum
from airflow.providers.standard.operators.bash import BashOperator
from airflow.sdk import DAG
with DAG(
dag_id="edge_adopt_demo",
schedule=None,
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
is_paused_upon_creation=False,
):
for i in range(1, 7):
BashOperator(task_id=f"task_{i}", bash_command="true")
EOF
uv run --frozen --project providers/edge3 airflow db migrate > /dev/null 2>&1
uv run --frozen --project providers/edge3 airflow dags reserialize >
/dev/null 2>&1
uv run --frozen --project providers/edge3 airflow dags trigger
edge_adopt_demo --run-id demo > /dev/null 2>&1
uv run --frozen --project providers/edge3 airflow scheduler --num-runs 5 >
"$AIRFLOW_HOME/scheduler-1.log" 2>&1
uv run --frozen --project providers/edge3 airflow tasks states-for-dag-run
edge_adopt_demo demo -o plain 2>/dev/null | grep -v FileType
uv run --frozen --project providers/edge3 airflow scheduler --num-runs 5 >
"$AIRFLOW_HOME/scheduler-2.log" 2>&1; echo "exit code: $?"
echo "DetachedInstanceError: $(grep -c "DetachedInstanceError:"
"$AIRFLOW_HOME/scheduler-2.log")"
echo "parallelism limit reached: $(grep -c "parallelism limit reached"
"$AIRFLOW_HOME/scheduler-2.log")"
```
The first `airflow scheduler --num-runs 5` queues three of the six tasks,
because `[core] parallelism` is 3, and then exits. The second one is the
restart. It should adopt the three task instances that the first scheduler left
queued.
On this branch, the restart adopts them. The adopted task instances keep
their slots, so the scheduler logs `Executor parallelism limit reached` in each
of its five loops and the other three tasks stay scheduled:
```text
dag_id logical_date task_id state start_date end_date
edge_adopt_demo task_1 queued
edge_adopt_demo task_2 queued
edge_adopt_demo task_3 queued
edge_adopt_demo task_4 scheduled
edge_adopt_demo task_5 scheduled
edge_adopt_demo task_6 scheduled
exit code: 0
DetachedInstanceError: 0
parallelism limit reached: 5
```
On `main`, the task states are the same, but the restart crashes before it
adopts anything:
```text
exit code: 1
DetachedInstanceError: 2
parallelism limit reached: 0
```
<!-- SPDX-License-Identifier: Apache-2.0
https://www.apache.org/licenses/LICENSE-2.0 -->
<!--
Thank you for contributing!
Please provide above a brief description of the changes made in this pull
request.
Write a good git commit message following this guide:
https://chris.beams.io/posts/git-commit/
Please make sure that your code changes are covered with tests.
And in case of new features or big changes remember to adjust the
documentation.
For user-facing UI changes, please attach before/after screenshots (or a
short
screen recording) so reviewers can assess the visual impact.
Feel free to ping (in general) for the review if you do not see reaction for
a few days
(72 Hours is the minimum reaction time you can expect from volunteers) - we
sometimes miss notifications.
In case of an existing issue, reference it using one of the following:
* closes: #ISSUE
* related: #ISSUE
-->
---
##### Was generative AI tooling used to co-author this PR?
<!--
If generative AI tooling has been used in the process of authoring this PR,
please
change below checkbox to `[X]` followed by the name of the tool, uncomment
the "Generated-by".
-->
- [X] Yes - Claude Code
<!--
Generated-by: [Tool Name] following [the
guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions)
-->
---
* Read the **[Pull Request
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#pull-request-guidelines)**
for more information. Note: commit author/co-author name and email in commits
become permanently public when merged.
* For fundamental code changes, an Airflow Improvement Proposal
([AIP](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
is needed.
* When adding dependency, check compliance with the [ASF 3rd Party License
Policy](https://www.apache.org/legal/resolved.html#category-x).
* For significant user-facing changes create newsfragment:
`{pr_number}.significant.rst`, in
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
You can add this file in a follow-up commit after the PR is created so you
know the PR number.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]