Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 28 additions & 23 deletions airflow-core/adr/lang-sdk/0013-persisted-task-handler-bindings.md
Original file line number Diff line number Diff line change
Expand Up @@ -271,33 +271,33 @@ none. The answer must not depend on the request, because the manager records it
every Dag file that resolves against the artifact.

**Dag-parsing child → manager.** [`DagFileParsingResult`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L133-L145) gains the resolved
bindings.
bindings and the probed artifacts.

```python
class SDKTaskHandlerBinding(BaseModel):
class TaskHandlerBinding(BaseModel):
dag_id: str
task_id: str
artifact_bundle_name: str
artifact_rel_path: str
artifact_size_bytes: int
artifact_cache_digest: str
handler_params: list[TaskHandlerParam]


class DagFileParsingResult(BaseModel):
fileloc: str
serialized_dags: list[LazyDeserializedDAG]
warnings: list | None = None
import_errors: dict[str, str] | None = None
task_handler_bindings: list[SDKTaskHandlerBinding] | None = None # new
task_handler_bindings: list[TaskHandlerBinding] | None = None # new
probed_artifacts: list[TaskHandlerArtifact] = [] # new
```

`None` and `[]` mean different things, and the difference is load-bearing:
`probed_artifacts` holds every artifact the parse probed, with its answer; it is recorded even when the bindings are `None`, and it is the only way an answer gets written.

For `task_handler_bindings`, `None` and `[]` mean different things, and the difference is load-bearing:

| value | meaning | manager does |
|--------|------------------------------------------|-----------------------|
| `None` | handlers were not evaluated in this parse | **nothing** — no reconcile |
| `[]` | evaluated, this file has no stub handlers | delete this file's rows |
| `[]` | evaluated, this file has no stub handlers | delete the rows of this result's Dags |
| `[…]` | evaluated, these are the bindings | reconcile to this set |

`None` covers the stability-check early return ([`_parse_file`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L245-L251)), callback-only runs
Expand Down Expand Up @@ -413,29 +413,34 @@ DagFileProcessorProcess(etl.py) [no DB — client con
on mismatch → import_errors[etl.py]=… AND bindings=None
│
├── DagFileParsingResult(serialized_dags=[…], import_errors={…},
▼ task_handler_bindings=[…] | [] | None)
DagProcessorManager.persist_parsing_result [writes DB]
└── update_dag_parsing_results_in_db — one transaction, in this order:
▼ task_handler_bindings=[…] | [] | None,
probed_artifacts=[…])
DagProcessorManager.handle_parsing_result [writes DB]
├── record the probed artifacts, a transaction of their own, skipped when there are none ← new
│ lang_sdk_task_handler_artifact UPSERT by (bundle_name, relative_fileloc):
│ fingerprint, answer, last_probed_at
│
└── persist_parsing_result → update_dag_parsing_results_in_db — one transaction, in this order:
1. add_dags / update_dags → DagModel rows exist
2. asset reference tables (existing)
3. lang_sdk_task_handler_artifact UPSERT by (bundle_name, relative_fileloc) ← new
4. lang_sdk_task_handler reconcile by dag_id ← new
5. SerializedDagModel / DagVersion / DagCode
6. ParseImportError / DagWarning
3. lang_sdk_task_handler reconcile by dag_id, artifact ids looked up by key ← new
4. SerializedDagModel / DagVersion / DagCode
5. ParseImportError / DagWarning
```

Step 3 must upsert: two parse children can discover the same artifact in the same loop and race on
the unique key.
The probed-artifact write must upsert: two parse children can probe the same artifact in the same loop and race on the unique key. It commits before the persist starts, so its exclusive row locks are gone before step 3 takes shared ones, and an answer is kept even when the persist fails.

Step 4 reconciles **by the `dag_id`s in the result**, not by file path: delete rows for those
Step 3 reconciles **by the `dag_id`s in the result**, not by file path: delete rows for those
`dag_id`s whose `task_id` is absent from the returned set, then insert or update the rest. Path-keyed
eviction breaks when a Dag moves between files — the old rows stay keyed to a path nothing parses any
more, and the primary key then blocks the new insert. With `dag_id` as the key a move simply updates
`dag_relative_fileloc`, and a Dag that disappears entirely is reclaimed by the `ON DELETE CASCADE`.
`dag_relative_fileloc`. A Dag removed from its file is only marked stale, so its rows stay until its `dag` row is deleted, and the `ON DELETE CASCADE` then removes them. Until then they keep their artifacts from the orphan sweep.

Step 3 writes nothing to the artifact table. A Dag bound to an artifact outside the Dag file's scope, or to one with no row because a concurrent orphan sweep deleted it, keeps its rows as they are, with a warning; in the second case the next parse finds no recorded answer and probes again.

Artifact rows are **never** evicted from one file's result. One artifact backs handlers owned by many
Python files, so this file seeing fewer candidates says nothing about another file's. They are a
cache; they are reclaimed by orphan sweep or by `db clean`, never by a per-file reconcile.
cache; they are reclaimed by orphan sweep, never by a per-file reconcile.

### Flow 2) Scheduling: the read path

Expand Down Expand Up @@ -557,8 +562,8 @@ argument error instead of catching it as an import error.

### Failure handling

**Validation fails.** Record the import error; leave every existing row untouched. `bindings=None`
already expresses "do not reconcile", so this needs no additional mechanism.
**Validation fails.** Record the import error; leave every binding row untouched. `bindings=None`
already expresses "do not reconcile", so this needs no additional mechanism. The answers probed for it are still recorded, so the next parse validates against them without probing.

**A stub task has no binding at all.** The scheduler fails it with the reason rather than queueing a
workload that would die on the worker at `ValueError("dag_path is required")`, far from the cause.
Expand All @@ -578,7 +583,7 @@ definition whose author can act) naming both artifact paths, since the fix is in
that owns rows in it; `lang_sdk_task_handler_artifact` is a cache with no per-file eviction.
- The artifact bundle is a second bundle on the execution path. `ExecuteTask` and `StartupDetails`
each grow one optional `SDKTaskHandlerRef`, and the worker performs a second `initialize()`. Workload payloads grow by roughly one `BundleInfo`.
- `DagFileParseRequest` and `DagFileParsingResult` each gain a field, and `ToSDKTaskHandlerProcessor`
- `DagFileParseRequest` gains a field and `DagFileParsingResult` two, and `ToSDKTaskHandlerProcessor`
becomes a fifth union the supervisor-schema registry introspects. Both messages already appear in
the generated schemas of all three SDKs, so the snapshot is regenerated and the two prek hooks guarding it run.
- Every Lang SDK runtime must answer `SDKTaskHandlerParseRequest`.
Expand Down
4 changes: 3 additions & 1 deletion airflow-core/docs/migrations-ref.rst
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@ Here's the list of all the Database Migrations that are executed via when you ru
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| Revision ID | Revises ID | Airflow Version | Description |
+=========================+==================+===================+==============================================================+
| ``f7ed13533d23`` (head) | ``90e4d18ccadf`` | ``3.4.0`` | Add lang_sdk_task_handler_artifact and lang_sdk_task_handler |
| ``37d645374a9c`` (head) | ``f7ed13533d23`` | ``3.4.0`` | Add last_probed_at index to lang_sdk_task_handler_artifact. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``f7ed13533d23`` | ``90e4d18ccadf`` | ``3.4.0`` | Add lang_sdk_task_handler_artifact and lang_sdk_task_handler |
| | | | tables. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``90e4d18ccadf`` | ``e5a91c7f42b3`` | ``3.4.0`` | Add timetable_asset_gated to DagModel. |
Expand Down
3 changes: 2 additions & 1 deletion airflow-core/src/airflow/config_templates/config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2727,7 +2727,8 @@ scheduler:
description: |
How often (in seconds) to check for stale DAGs (DAGs which are no longer present in
the expected files) which should be deactivated, as well as assets that are no longer
referenced and should be marked as orphaned.
referenced and should be marked as orphaned. The same check deletes the Lang-SDK task
handler artifacts that no stub task uses any more.
version_added: 2.5.0
type: integer
example: ~
Expand Down
Loading
Loading